Browse Source

filter nodes && added TbMsgType enum

pull/8786/head
ShvaykaD 3 years ago
parent
commit
22874e8a65
  1. 26
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  2. 4
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
  3. 4
      application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java
  4. 94
      application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java
  5. 25
      application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java
  6. 24
      application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java
  7. 10
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  8. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java
  9. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java
  10. 5
      application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java
  11. 5
      application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java
  12. 18
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  13. 3
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  14. 44
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  15. 8
      common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java
  16. 88
      common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java
  17. 95
      common/data/src/main/java/org/thingsboard/server/common/data/msg/TbMsgType.java
  18. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/EmptyNodeConfiguration.java
  19. 9
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNodeConnectionType.java
  20. 15
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java
  21. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java
  22. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCopyAttributesToEntityViewNode.java
  23. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbMsgCountNode.java
  24. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java
  25. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java
  26. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java
  27. 52
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java
  28. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbAssetTypeSwitchNode.java
  29. 37
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNode.java
  30. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNodeConfig.java
  31. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckMessageNode.java
  32. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNode.java
  33. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbDeviceTypeSwitchNode.java
  34. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java
  35. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java
  36. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java
  37. 86
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java
  38. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeFilterNode.java
  39. 50
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java
  40. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbCheckpointNode.java
  41. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/geo/TbGpsGeofencingFilterNode.java
  42. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java
  43. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java
  44. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java
  45. 31
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java
  46. 8
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java
  47. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java
  48. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
  49. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java
  50. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java
  51. 13
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateRelationNodeTest.java
  52. 18
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNodeTest.java
  53. 14
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsFilterNodeTest.java
  54. 15
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/DeviceStateTest.java
  55. 12
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java

26
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.actors.ruleChain;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import io.netty.channel.EventLoopGroup;
@ -34,7 +33,7 @@ import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.rule.engine.api.SmsService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.slack.SlackService;
import org.thingsboard.rule.engine.api.sms.SmsSenderFactory;
import org.thingsboard.rule.engine.util.TenantIdLoader;
@ -42,7 +41,6 @@ import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.Customer;
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;
@ -114,6 +112,10 @@ import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
/**
* Created by ashvayka on 19.03.18.
*/
@ -132,7 +134,7 @@ class DefaultTbContext implements TbContext {
@Override
public void tellSuccess(TbMsg msg) {
tellNext(msg, Collections.singleton(TbRelationTypes.SUCCESS), null);
tellNext(msg, Collections.singleton(TbNodeConnectionType.SUCCESS), null);
}
@Override
@ -211,7 +213,7 @@ class DefaultTbContext implements TbContext {
@Override
public void enqueueForTellFailure(TbMsg tbMsg, String failureMessage) {
TopicPartitionInfo tpi = resolvePartition(tbMsg);
enqueueForTellNext(tpi, tbMsg, Collections.singleton(TbRelationTypes.FAILURE), failureMessage, null, null);
enqueueForTellNext(tpi, tbMsg, Collections.singleton(TbNodeConnectionType.FAILURE), failureMessage, null, null);
}
@Override
@ -309,7 +311,7 @@ class DefaultTbContext implements TbContext {
@Override
public void tellFailure(TbMsg msg, Throwable th) {
if (nodeCtx.getSelf().isDebugMode()) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, TbRelationTypes.FAILURE, th);
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, TbNodeConnectionType.FAILURE, th);
}
String failureMessage;
if (th != null) {
@ -322,7 +324,7 @@ class DefaultTbContext implements TbContext {
failureMessage = null;
}
nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getRuleChainId(),
nodeCtx.getSelf().getId(), Collections.singleton(TbRelationTypes.FAILURE),
nodeCtx.getSelf().getId(), Collections.singleton(TbNodeConnectionType.FAILURE),
msg, failureMessage));
}
@ -346,7 +348,7 @@ class DefaultTbContext implements TbContext {
}
public TbMsg customerCreatedMsg(Customer customer, RuleNodeId ruleNodeId) {
return entityActionMsg(customer, customer.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
return entityActionMsg(customer, customer.getId(), ruleNodeId, ENTITY_CREATED.name());
}
public TbMsg deviceCreatedMsg(Device device, RuleNodeId ruleNodeId) {
@ -354,7 +356,7 @@ class DefaultTbContext implements TbContext {
if (device.getDeviceProfileId() != null) {
deviceProfile = mainCtx.getDeviceProfileCache().find(device.getDeviceProfileId());
}
return entityActionMsg(device, device.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, deviceProfile);
return entityActionMsg(device, device.getId(), ruleNodeId, ENTITY_CREATED.name(), deviceProfile);
}
public TbMsg assetCreatedMsg(Asset asset, RuleNodeId ruleNodeId) {
@ -362,7 +364,7 @@ class DefaultTbContext implements TbContext {
if (asset.getAssetProfileId() != null) {
assetProfile = mainCtx.getAssetProfileCache().find(asset.getAssetProfileId());
}
return entityActionMsg(asset, asset.getId(), ruleNodeId, DataConstants.ENTITY_CREATED, assetProfile);
return entityActionMsg(asset, asset.getId(), ruleNodeId, ENTITY_CREATED.name(), assetProfile);
}
public TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action) {
@ -382,7 +384,7 @@ class DefaultTbContext implements TbContext {
if (attributes != null) {
attributes.forEach(attributeKvEntry -> JacksonUtil.addKvEntry(entityNode, attributeKvEntry));
}
return attributesActionMsg(originator, ruleNodeId, scope, DataConstants.ATTRIBUTES_UPDATED, JacksonUtil.toString(entityNode));
return attributesActionMsg(originator, ruleNodeId, scope, ATTRIBUTES_UPDATED.name(), JacksonUtil.toString(entityNode));
}
public TbMsg attributesDeletedActionMsg(EntityId originator, RuleNodeId ruleNodeId, String scope, List<String> keys) {
@ -391,7 +393,7 @@ class DefaultTbContext implements TbContext {
if (keys != null) {
keys.forEach(attrsArrayNode::add);
}
return attributesActionMsg(originator, ruleNodeId, scope, DataConstants.ATTRIBUTES_DELETED, JacksonUtil.toString(entityNode));
return attributesActionMsg(originator, ruleNodeId, scope, ATTRIBUTES_DELETED.name(), JacksonUtil.toString(entityNode));
}
private TbMsg attributesActionMsg(EntityId originator, RuleNodeId ruleNodeId, String scope, String action, String msgData) {

4
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java

@ -16,7 +16,7 @@
package org.thingsboard.server.actors.ruleChain;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorCtx;
import org.thingsboard.server.actors.TbActorRef;
@ -307,7 +307,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
int relationsCount = relationsByTypes.size();
if (relationsCount == 0) {
log.trace("[{}][{}][{}] No outbound relations to process", tenantId, entityId, msg.getId());
if (relationTypes.contains(TbRelationTypes.FAILURE)) {
if (relationTypes.contains(TbNodeConnectionType.FAILURE)) {
RuleNodeCtx ruleNodeCtx = nodeActors.get(originatorNodeId);
if (ruleNodeCtx != null) {
msg.getCallback().onFailure(new RuleNodeException(failureMessage, ruleChainName, ruleNodeCtx.getSelf()));

4
application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java

@ -52,7 +52,7 @@ import org.thingsboard.server.service.security.permission.Operation;
import javax.annotation.Nullable;
import java.util.UUID;
import static org.thingsboard.server.common.data.DataConstants.RPC_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.RPC_DELETED;
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID;
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID_PARAM_DESCRIPTION;
import static org.thingsboard.server.controller.ControllerConstants.MARKDOWN_CODE_BLOCK_END;
@ -239,7 +239,7 @@ public class RpcV2Controller extends AbstractRpcController {
rpcService.deleteRpc(getTenantId(), rpcId);
rpc.setStatus(RpcStatus.DELETED);
TbMsg msg = TbMsg.newMsg(RPC_DELETED, rpc.getDeviceId(), TbMsgMetaData.EMPTY, JacksonUtil.toString(rpc));
TbMsg msg = TbMsg.newMsg(RPC_DELETED.name(), rpc.getDeviceId(), TbMsgMetaData.EMPTY, JacksonUtil.toString(rpc));
tbClusterService.pushMsgToRuleEngine(getTenantId(), rpc.getDeviceId(), msg, null);
}
}

94
application/src/main/java/org/thingsboard/server/service/action/EntityActionService.java

@ -26,7 +26,6 @@ import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmComment;
@ -38,19 +37,21 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.AlarmCommentTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.EntitiesLimitTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.EntityActionTrigger;
import org.thingsboard.server.common.data.relation.EntityRelation;
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.notification.NotificationRuleProcessor;
import org.thingsboard.server.common.data.notification.rule.trigger.AlarmAssignmentTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.AlarmCommentTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.EntitiesLimitTrigger;
import org.thingsboard.server.common.data.notification.rule.trigger.EntityActionTrigger;
import org.thingsboard.server.dao.audit.AuditLogService;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
@Service
@ -63,85 +64,8 @@ public class EntityActionService {
public void pushEntityActionToRuleEngine(EntityId entityId, HasName entity, TenantId tenantId, CustomerId customerId,
ActionType actionType, User user, Object... additionalInfo) {
String msgType = null;
switch (actionType) {
case ADDED:
msgType = DataConstants.ENTITY_CREATED;
break;
case DELETED:
msgType = DataConstants.ENTITY_DELETED;
break;
case UPDATED:
msgType = DataConstants.ENTITY_UPDATED;
break;
case ASSIGNED_TO_CUSTOMER:
msgType = DataConstants.ENTITY_ASSIGNED;
break;
case UNASSIGNED_FROM_CUSTOMER:
msgType = DataConstants.ENTITY_UNASSIGNED;
break;
case ATTRIBUTES_UPDATED:
msgType = DataConstants.ATTRIBUTES_UPDATED;
break;
case ATTRIBUTES_DELETED:
msgType = DataConstants.ATTRIBUTES_DELETED;
break;
case ALARM_ACK:
msgType = DataConstants.ALARM_ACK;
break;
case ALARM_CLEAR:
msgType = DataConstants.ALARM_CLEAR;
break;
case ALARM_ASSIGNED:
msgType = DataConstants.ALARM_ASSIGNED;
break;
case ALARM_UNASSIGNED:
msgType = DataConstants.ALARM_UNASSIGNED;
break;
case ALARM_DELETE:
msgType = DataConstants.ALARM_DELETE;
break;
case ADDED_COMMENT:
msgType = DataConstants.COMMENT_CREATED;
break;
case UPDATED_COMMENT:
msgType = DataConstants.COMMENT_UPDATED;
break;
case ASSIGNED_FROM_TENANT:
msgType = DataConstants.ENTITY_ASSIGNED_FROM_TENANT;
break;
case ASSIGNED_TO_TENANT:
msgType = DataConstants.ENTITY_ASSIGNED_TO_TENANT;
break;
case PROVISION_SUCCESS:
msgType = DataConstants.PROVISION_SUCCESS;
break;
case PROVISION_FAILURE:
msgType = DataConstants.PROVISION_FAILURE;
break;
case TIMESERIES_UPDATED:
msgType = DataConstants.TIMESERIES_UPDATED;
break;
case TIMESERIES_DELETED:
msgType = DataConstants.TIMESERIES_DELETED;
break;
case ASSIGNED_TO_EDGE:
msgType = DataConstants.ENTITY_ASSIGNED_TO_EDGE;
break;
case UNASSIGNED_FROM_EDGE:
msgType = DataConstants.ENTITY_UNASSIGNED_FROM_EDGE;
break;
case RELATION_ADD_OR_UPDATE:
msgType = DataConstants.RELATION_ADD_OR_UPDATE;
break;
case RELATION_DELETED:
msgType = DataConstants.RELATION_DELETED;
break;
case RELATIONS_DELETED:
msgType = DataConstants.RELATIONS_DELETED;
break;
}
if (!StringUtils.isEmpty(msgType)) {
Optional<TbMsgType> msgType = actionType.getRuleEngineMsgType();
if (msgType.isPresent()) {
try {
TbMsgMetaData metaData = new TbMsgMetaData();
if (user != null) {
@ -247,7 +171,7 @@ public class EntityActionService {
if (tenantId != null && !tenantId.isSysTenantId()) {
processNotificationRules(tenantId, entityId, entity, actionType, user, additionalInfo);
}
TbMsg tbMsg = TbMsg.newMsg(msgType, entityId, customerId, metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
TbMsg tbMsg = TbMsg.newMsg(msgType.get().name(), entityId, customerId, metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
tbClusterService.pushMsgToRuleEngine(tenantId, entityId, tbMsg, null);
} catch (Exception e) {
log.warn("[{}] Failed to push entity action to rule engine: {}", entityId, actionType, e);

25
application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java

@ -30,8 +30,12 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.rule.engine.api.NodeDefinition;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.TbVersionedNode;
import org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode;
import org.thingsboard.rule.engine.filter.TbOriginatorTypeSwitchNode;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentDescriptor;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -194,7 +198,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
scannedComponent.setName(ruleNodeAnnotation.name());
scannedComponent.setScope(ruleNodeAnnotation.scope());
scannedComponent.setClusteringMode(ruleNodeAnnotation.clusteringMode());
NodeDefinition nodeDefinition = prepareNodeDefinition(ruleNodeAnnotation);
NodeDefinition nodeDefinition = prepareNodeDefinition(clazz, ruleNodeAnnotation);
ObjectNode configurationDescriptor = JacksonUtil.newObjectNode();
JsonNode node = JacksonUtil.valueToTree(nodeDefinition);
configurationDescriptor.set("nodeDefinition", node);
@ -221,13 +225,13 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
return scannedComponent;
}
private NodeDefinition prepareNodeDefinition(RuleNode nodeAnnotation) throws Exception {
private NodeDefinition prepareNodeDefinition(Class<?> clazz, RuleNode nodeAnnotation) throws Exception {
NodeDefinition nodeDefinition = new NodeDefinition();
nodeDefinition.setDetails(nodeAnnotation.nodeDetails());
nodeDefinition.setDescription(nodeAnnotation.nodeDescription());
nodeDefinition.setInEnabled(nodeAnnotation.inEnabled());
nodeDefinition.setOutEnabled(nodeAnnotation.outEnabled());
nodeDefinition.setRelationTypes(getRelationTypesWithFailureRelation(nodeAnnotation));
nodeDefinition.setRelationTypes(getRelationTypesWithFailureRelation(clazz, nodeAnnotation));
nodeDefinition.setCustomRelations(nodeAnnotation.customRelations());
nodeDefinition.setRuleChainNode(nodeAnnotation.ruleChainNode());
Class<? extends NodeConfiguration> configClazz = nodeAnnotation.configClazz();
@ -242,10 +246,17 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
return nodeDefinition;
}
private String[] getRelationTypesWithFailureRelation(RuleNode nodeAnnotation) {
private String[] getRelationTypesWithFailureRelation(Class<?> clazz, RuleNode nodeAnnotation) {
List<String> relationTypes = new ArrayList<>(Arrays.asList(nodeAnnotation.relationTypes()));
if (!relationTypes.contains(TbRelationTypes.FAILURE)) {
relationTypes.add(TbRelationTypes.FAILURE);
if (TbOriginatorTypeSwitchNode.class.equals(clazz)) {
relationTypes.addAll(EntityType.NORMAL_NAMES);
}
if (TbMsgTypeSwitchNode.class.equals(clazz)) {
relationTypes.addAll(TbMsgType.NODE_CONNECTIONS);
relationTypes.add(TbMsgType.OTHER);
}
if (!relationTypes.contains(TbNodeConnectionType.FAILURE)) {
relationTypes.add(TbNodeConnectionType.FAILURE);
}
return relationTypes.toArray(new String[relationTypes.size()]);
}

24
application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java

@ -23,7 +23,6 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.cluster.TbClusterService;
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.DeviceProfileProvisionType;
@ -69,6 +68,11 @@ import java.util.concurrent.ExecutionException;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.PROVISION_FAILURE;
import static org.thingsboard.server.common.data.msg.TbMsgType.PROVISION_SUCCESS;
@Service
@Slf4j
@ -162,7 +166,7 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
if (targetProfile.getProfileData().getProvisionConfiguration().getProvisionDeviceSecret().equals(provisionRequestSecret)) {
if (targetDevice != null) {
log.warn("[{}] The device is present and could not be provisioned once more!", targetDevice.getName());
notify(targetDevice, provisionRequest, DataConstants.PROVISION_FAILURE, false);
notify(targetDevice, provisionRequest, PROVISION_FAILURE.name(), false);
throw new ProvisionFailedException(ProvisionResponseStatus.FAILURE.name());
} else {
return createDevice(provisionRequest, targetProfile);
@ -188,13 +192,13 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
private ProvisionResponse processProvision(Device device, ProvisionRequest provisionRequest) {
try {
Optional<AttributeKvEntry> provisionState = attributesService.find(device.getTenantId(), device.getId(),
DataConstants.SERVER_SCOPE, DEVICE_PROVISION_STATE).get();
SERVER_SCOPE, DEVICE_PROVISION_STATE).get();
if (provisionState != null && provisionState.isPresent() && !provisionState.get().getValueAsString().equals(PROVISIONED_STATE)) {
notify(device, provisionRequest, DataConstants.PROVISION_FAILURE, false);
notify(device, provisionRequest, PROVISION_FAILURE.name(), false);
throw new ProvisionFailedException(ProvisionResponseStatus.FAILURE.name());
} else {
saveProvisionStateAttribute(device).get();
notify(device, provisionRequest, DataConstants.PROVISION_SUCCESS, true);
notify(device, provisionRequest, PROVISION_SUCCESS.name(), true);
}
} catch (InterruptedException | ExecutionException e) {
throw new ProvisionFailedException(ProvisionResponseStatus.FAILURE.name());
@ -222,14 +226,14 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
clusterService.onDeviceUpdated(savedDevice, null);
saveProvisionStateAttribute(savedDevice).get();
pushDeviceCreatedEventToRuleEngine(savedDevice);
notify(savedDevice, provisionRequest, DataConstants.PROVISION_SUCCESS, true);
notify(savedDevice, provisionRequest, PROVISION_SUCCESS.name(), true);
return new ProvisionResponse(getDeviceCredentials(savedDevice), ProvisionResponseStatus.SUCCESS);
} catch (Exception e) {
log.warn("[{}] Error during device creation from provision request: [{}]", provisionRequest.getDeviceName(), provisionRequest, e);
Device device = deviceService.findDeviceByTenantIdAndName(profile.getTenantId(), provisionRequest.getDeviceName());
if (device != null) {
notify(device, provisionRequest, DataConstants.PROVISION_FAILURE, false);
notify(device, provisionRequest, PROVISION_FAILURE.name(), false);
}
throw new ProvisionFailedException(ProvisionResponseStatus.FAILURE.name());
}
@ -244,7 +248,7 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
}
private ListenableFuture<List<String>> saveProvisionStateAttribute(Device device) {
return attributesService.save(device.getTenantId(), device.getId(), DataConstants.SERVER_SCOPE,
return attributesService.save(device.getTenantId(), device.getId(), SERVER_SCOPE,
Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry(DEVICE_PROVISION_STATE, PROVISIONED_STATE),
System.currentTimeMillis())));
}
@ -266,10 +270,10 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
private void pushDeviceCreatedEventToRuleEngine(Device device) {
try {
ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.valueToTree(device);
TbMsg msg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, device.getId(), device.getCustomerId(), createTbMsgMetaData(device), JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode));
TbMsg msg = TbMsg.newMsg(ENTITY_CREATED.name(), device.getId(), device.getCustomerId(), createTbMsgMetaData(device), JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode));
sendToRuleEngine(device.getTenantId(), msg, null);
} catch (JsonProcessingException | IllegalArgumentException e) {
log.warn("[{}] Failed to push device action to rule engine: {}", device.getId(), DataConstants.ENTITY_CREATED, e);
log.warn("[{}] Failed to push device action to rule engine: {}", device.getId(), ENTITY_CREATED, e);
}
}

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

@ -71,9 +71,9 @@ import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
import static org.thingsboard.server.common.data.DataConstants.CONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
import static org.thingsboard.server.common.data.msg.TbMsgType.CONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.DISCONNECT_EVENT;
@Service
@Slf4j
@ -278,7 +278,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
long lastConnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, lastConnectTs);
pushRuleEngineMessage(edgeGrpcSession.getEdge().getTenantId(), edgeId, lastConnectTs, CONNECT_EVENT);
pushRuleEngineMessage(edgeGrpcSession.getEdge().getTenantId(), edgeId, lastConnectTs, CONNECT_EVENT.name());
cancelScheduleEdgeEventsCheck(edgeId);
scheduleEdgeEventsCheck(edgeGrpcSession);
}
@ -395,7 +395,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
long lastDisconnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, lastDisconnectTs);
pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edgeId, lastDisconnectTs, DISCONNECT_EVENT);
pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edgeId, lastDisconnectTs, DISCONNECT_EVENT.name());
cancelScheduleEdgeEventsCheck(edgeId);
} else {
log.debug("[{}] edge session [{}] is not available anymore, nothing to remove. most probably this session is already outdated!", edgeId, sessionId);
@ -451,7 +451,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void pushRuleEngineMessage(TenantId tenantId, EdgeId edgeId, long ts, String msgType) {
try {
ObjectNode edgeState = JacksonUtil.newObjectNode();
if (msgType.equals(CONNECT_EVENT)) {
if (msgType.equals(CONNECT_EVENT.name())) {
edgeState.put(DefaultDeviceStateService.ACTIVITY_STATE, true);
edgeState.put(DefaultDeviceStateService.LAST_CONNECT_TIME, ts);
} else {

6
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java

@ -62,6 +62,8 @@ import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg;
import java.util.UUID;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
@Component
@Slf4j
@TbCoreComponent
@ -124,7 +126,7 @@ public class DeviceEdgeProcessor extends BaseDeviceProcessor {
try {
Device device = deviceService.findDeviceById(tenantId, deviceId);
ObjectNode entityNode = JacksonUtil.OBJECT_MAPPER.valueToTree(device);
TbMsg tbMsg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, device.getCustomerId(),
TbMsg tbMsg = TbMsg.newMsg(ENTITY_CREATED.name(), deviceId, device.getCustomerId(),
getActionTbMsgMetaData(edge, device.getCustomerId()), TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(entityNode));
tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() {
@Override
@ -138,7 +140,7 @@ public class DeviceEdgeProcessor extends BaseDeviceProcessor {
}
});
} catch (JsonProcessingException | IllegalArgumentException e) {
log.warn("[{}] Failed to push device action to rule engine: {}", deviceId, DataConstants.ENTITY_CREATED, e);
log.warn("[{}] Failed to push device action to rule engine: {}", deviceId, ENTITY_CREATED.name(), e);
}
}

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java

@ -73,6 +73,8 @@ import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
@Slf4j
public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@ -257,7 +259,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override
public void onSuccess(@Nullable Void tmp) {
var defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId);
TbMsg tbMsg = TbMsg.newMsg(defaultQueueAndRuleChain.getKey(), DataConstants.ATTRIBUTES_UPDATED, entityId,
TbMsg tbMsg = TbMsg.newMsg(defaultQueueAndRuleChain.getKey(), ATTRIBUTES_UPDATED.name(), entityId,
customerId, metaData, gson.toJson(json), defaultQueueAndRuleChain.getValue(), null);
tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() {
@Override

5
application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java

@ -21,7 +21,6 @@ import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasName;
@ -53,6 +52,8 @@ import org.thingsboard.server.service.gateway_device.GatewayNotificationsService
import java.util.List;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_ASSIGNED_FROM_TENANT;
@Slf4j
@Service
@RequiredArgsConstructor
@ -286,7 +287,7 @@ public class DefaultTbNotificationEntityService implements TbNotificationEntityS
private void pushAssignedFromNotification(Tenant currentTenant, TenantId newTenantId, Device assignedDevice) {
String data = JacksonUtil.toString(JacksonUtil.valueToTree(assignedDevice));
if (data != null) {
TbMsg tbMsg = TbMsg.newMsg(DataConstants.ENTITY_ASSIGNED_FROM_TENANT, assignedDevice.getId(),
TbMsg tbMsg = TbMsg.newMsg(ENTITY_ASSIGNED_FROM_TENANT.name(), assignedDevice.getId(),
assignedDevice.getCustomerId(), getMetaDataForAssignedFrom(currentTenant), TbMsgDataType.JSON, data);
tbClusterService.pushMsgToRuleEngine(newTenantId, assignedDevice.getId(), tbMsg, null);
}

5
application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.service.rpc;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -49,6 +48,8 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import static org.thingsboard.server.common.data.msg.TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE;
/**
* Created by ashvayka on 27.03.18.
*/
@ -182,7 +183,7 @@ public class DefaultTbCoreDeviceRpcService implements TbCoreDeviceRpcService {
entityNode.put(DataConstants.ADDITIONAL_INFO, msg.getAdditionalInfo());
try {
TbMsg tbMsg = TbMsg.newMsg(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE, msg.getDeviceId(), Optional.ofNullable(currentUser).map(User::getCustomerId).orElse(null), metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
TbMsg tbMsg = TbMsg.newMsg(RPC_CALL_FROM_SERVER_TO_DEVICE.name(), msg.getDeviceId(), Optional.ofNullable(currentUser).map(User::getCustomerId).orElse(null), metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
clusterService.pushMsgToRuleEngine(msg.getTenantId(), msg.getDeviceId(), tbMsg, null);
} catch (IllegalArgumentException e) {
throw new RuntimeException(e);

18
application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

@ -51,6 +51,7 @@ import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityTrigger;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.data.query.EntityData;
@ -63,7 +64,6 @@ 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.notification.NotificationRuleProcessor;
import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityTrigger;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -102,11 +102,11 @@ import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.DataConstants.CONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.INACTIVITY_EVENT;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.CONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT;
/**
* Created by ashvayka on 01.05.18.
@ -229,7 +229,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
long ts = System.currentTimeMillis();
stateData.getState().setLastConnectTime(ts);
save(deviceId, LAST_CONNECT_TIME, ts);
pushRuleEngineMessage(stateData, CONNECT_EVENT);
pushRuleEngineMessage(stateData, CONNECT_EVENT.name());
checkAndUpdateState(deviceId, stateData);
}
@ -271,7 +271,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
long ts = System.currentTimeMillis();
stateData.getState().setLastDisconnectTime(ts);
save(deviceId, LAST_DISCONNECT_TIME, ts);
pushRuleEngineMessage(stateData, DISCONNECT_EVENT);
pushRuleEngineMessage(stateData, DISCONNECT_EVENT.name());
}
@Override
@ -533,7 +533,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
private void onDeviceActivityStatusChange(DeviceId deviceId, boolean active, DeviceStateData stateData) {
save(deviceId, ACTIVITY_STATE, active);
pushRuleEngineMessage(stateData, active ? ACTIVITY_EVENT : INACTIVITY_EVENT);
pushRuleEngineMessage(stateData, active ? ACTIVITY_EVENT.name() : INACTIVITY_EVENT.name());
TbMsgMetaData metaData = stateData.getMetaData();
notificationRuleProcessor.process(DeviceActivityTrigger.builder()
.tenantId(stateData.getTenantId()).customerId(stateData.getCustomerId())
@ -775,7 +775,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
DeviceState state = stateData.getState();
try {
String data;
if (msgType.equals(CONNECT_EVENT)) {
if (msgType.equals(CONNECT_EVENT.name())) {
ObjectNode stateNode = JacksonUtil.convertValue(state, ObjectNode.class);
stateNode.remove(ACTIVITY_STATE);
data = JacksonUtil.toString(stateNode);

3
application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java

@ -113,6 +113,7 @@ import java.util.concurrent.locks.ReentrantLock;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.PASSWORD_MISMATCH;
import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.VALID;
@ -346,7 +347,7 @@ public class DefaultTransportApiService implements TransportApiService {
DeviceId deviceId = device.getId();
JsonNode entityNode = JacksonUtil.valueToTree(device);
TbMsg tbMsg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, customerId, metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
TbMsg tbMsg = TbMsg.newMsg(ENTITY_CREATED.name(), deviceId, customerId, metaData, TbMsgDataType.JSON, JacksonUtil.toString(entityNode));
tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, null);
} else {
JsonNode deviceAdditionalInfo = device.getAdditionalInfo();

44
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -54,53 +54,9 @@ public class DataConstants {
return new String[]{CLIENT_SCOPE, SHARED_SCOPE, SERVER_SCOPE};
}
public static final String ALARM = "ALARM";
public static final String IN = "IN";
public static final String OUT = "OUT";
public static final String INACTIVITY_EVENT = "INACTIVITY_EVENT";
public static final String CONNECT_EVENT = "CONNECT_EVENT";
public static final String DISCONNECT_EVENT = "DISCONNECT_EVENT";
public static final String ACTIVITY_EVENT = "ACTIVITY_EVENT";
public static final String ENTITY_CREATED = "ENTITY_CREATED";
public static final String ENTITY_UPDATED = "ENTITY_UPDATED";
public static final String ENTITY_DELETED = "ENTITY_DELETED";
public static final String ENTITY_ASSIGNED = "ENTITY_ASSIGNED";
public static final String ENTITY_UNASSIGNED = "ENTITY_UNASSIGNED";
public static final String ATTRIBUTES_UPDATED = "ATTRIBUTES_UPDATED";
public static final String ATTRIBUTES_DELETED = "ATTRIBUTES_DELETED";
public static final String TIMESERIES_UPDATED = "TIMESERIES_UPDATED";
public static final String TIMESERIES_DELETED = "TIMESERIES_DELETED";
public static final String ALARM_ACK = "ALARM_ACK";
public static final String ALARM_CLEAR = "ALARM_CLEAR";
public static final String ALARM_ASSIGNED = "ALARM_ASSIGNED";
public static final String ALARM_UNASSIGNED = "ALARM_UNASSIGNED";
public static final String ALARM_DELETE = "ALARM_DELETE";
public static final String COMMENT_CREATED = "COMMENT_CREATED";
public static final String COMMENT_UPDATED = "COMMENT_UPDATED";
public static final String ENTITY_ASSIGNED_FROM_TENANT = "ENTITY_ASSIGNED_FROM_TENANT";
public static final String ENTITY_ASSIGNED_TO_TENANT = "ENTITY_ASSIGNED_TO_TENANT";
public static final String PROVISION_SUCCESS = "PROVISION_SUCCESS";
public static final String PROVISION_FAILURE = "PROVISION_FAILURE";
public static final String ENTITY_ASSIGNED_TO_EDGE = "ENTITY_ASSIGNED_TO_EDGE";
public static final String ENTITY_UNASSIGNED_FROM_EDGE = "ENTITY_UNASSIGNED_FROM_EDGE";
public static final String RELATION_ADD_OR_UPDATE = "RELATION_ADD_OR_UPDATE";
public static final String RELATION_DELETED = "RELATION_DELETED";
public static final String RELATIONS_DELETED = "RELATIONS_DELETED";
public static final String RPC_CALL_FROM_SERVER_TO_DEVICE = "RPC_CALL_FROM_SERVER_TO_DEVICE";
public static final String RPC_QUEUED = "RPC_QUEUED";
public static final String RPC_SENT = "RPC_SENT";
public static final String RPC_DELIVERED = "RPC_DELIVERED";
public static final String RPC_SUCCESSFUL = "RPC_SUCCESSFUL";
public static final String RPC_TIMEOUT = "RPC_TIMEOUT";
public static final String RPC_EXPIRED = "RPC_EXPIRED";
public static final String RPC_FAILED = "RPC_FAILED";
public static final String RPC_DELETED = "RPC_DELETED";
public static final String DEFAULT_SECRET_KEY = "";
public static final String SECRET_KEY_FIELD_NAME = "secretKey";
public static final String DURATION_MS_FIELD_NAME = "durationMs";

8
common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java

@ -18,6 +18,10 @@ package org.thingsboard.server.common.data;
import lombok.Getter;
import org.apache.commons.lang3.StringUtils;
import java.util.EnumSet;
import java.util.List;
import java.util.stream.Collectors;
/**
* @author Andrew Shvayka
*/
@ -49,6 +53,10 @@ public enum EntityType {
NOTIFICATION,
NOTIFICATION_RULE;
public static final List<String> NORMAL_NAMES = EnumSet.allOf(EntityType.class).stream()
.map(EntityType::getNormalName).collect(Collectors.toUnmodifiableList());
@Getter
private final String normalName = StringUtils.capitalize(StringUtils.removeStart(name(), "TB_")
.toLowerCase().replaceAll("_", " "));

88
common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java

@ -16,49 +16,61 @@
package org.thingsboard.server.common.data.audit;
import lombok.Getter;
import org.thingsboard.server.common.data.msg.TbMsgType;
import java.util.Optional;
@Getter
public enum ActionType {
ADDED(false), // log entity
DELETED(false), // log string id
UPDATED(false), // log entity
ATTRIBUTES_UPDATED(false), // log attributes/values
ATTRIBUTES_DELETED(false), // log attributes
TIMESERIES_UPDATED(false), // log timeseries update
TIMESERIES_DELETED(false), // log timeseries
RPC_CALL(false), // log method and params
CREDENTIALS_UPDATED(false), // log new credentials
ASSIGNED_TO_CUSTOMER(false), // log customer name
UNASSIGNED_FROM_CUSTOMER(false), // log customer name
ACTIVATED(false), // log string id
SUSPENDED(false), // log string id
CREDENTIALS_READ(true), // log device id
ATTRIBUTES_READ(true), // log attributes
RELATION_ADD_OR_UPDATE(false),
RELATION_DELETED(false),
RELATIONS_DELETED(false),
ALARM_ACK(false),
ALARM_CLEAR(false),
ALARM_DELETE(false),
ALARM_ASSIGNED(false),
ALARM_UNASSIGNED(false),
LOGIN(false),
LOGOUT(false),
LOCKOUT(false),
ASSIGNED_FROM_TENANT(false),
ASSIGNED_TO_TENANT(false),
PROVISION_SUCCESS(false),
PROVISION_FAILURE(false),
ASSIGNED_TO_EDGE(false), // log edge name
UNASSIGNED_FROM_EDGE(false),
ADDED_COMMENT(false),
UPDATED_COMMENT(false),
DELETED_COMMENT(false),
SMS_SENT(false);
ADDED(false, TbMsgType.ENTITY_CREATED), // log entity
DELETED(false, TbMsgType.ENTITY_DELETED), // log string id
UPDATED(false, TbMsgType.ENTITY_UPDATED), // log entity
ATTRIBUTES_UPDATED(false, TbMsgType.ATTRIBUTES_UPDATED), // log attributes/values
ATTRIBUTES_DELETED(false, TbMsgType.ATTRIBUTES_DELETED), // log attributes
TIMESERIES_UPDATED(false, TbMsgType.TIMESERIES_UPDATED), // log timeseries update
TIMESERIES_DELETED(false, TbMsgType.TIMESERIES_DELETED), // log timeseries
RPC_CALL(false, null), // log method and params
CREDENTIALS_UPDATED(false, null), // log new credentials
ASSIGNED_TO_CUSTOMER(false, TbMsgType.ENTITY_ASSIGNED), // log customer name
UNASSIGNED_FROM_CUSTOMER(false, TbMsgType.ENTITY_UNASSIGNED), // log customer name
ACTIVATED(false, null), // log string id
SUSPENDED(false, null), // log string id
CREDENTIALS_READ(true, null), // log device id
ATTRIBUTES_READ(true, null), // log attributes
RELATION_ADD_OR_UPDATE(false, TbMsgType.RELATION_ADD_OR_UPDATE),
RELATION_DELETED(false, TbMsgType.RELATION_DELETED),
RELATIONS_DELETED(false, TbMsgType.RELATIONS_DELETED),
ALARM_ACK(false, TbMsgType.ALARM_ACK),
ALARM_CLEAR(false, TbMsgType.ALARM_CLEAR),
ALARM_DELETE(false, TbMsgType.ALARM_DELETE),
ALARM_ASSIGNED(false, TbMsgType.ALARM_ASSIGNED),
ALARM_UNASSIGNED(false, TbMsgType.ALARM_UNASSIGNED),
LOGIN(false, null),
LOGOUT(false, null),
LOCKOUT(false, null),
ASSIGNED_FROM_TENANT(false, TbMsgType.ENTITY_ASSIGNED_FROM_TENANT),
ASSIGNED_TO_TENANT(false, TbMsgType.ENTITY_ASSIGNED_TO_TENANT),
PROVISION_SUCCESS(false, TbMsgType.PROVISION_SUCCESS),
PROVISION_FAILURE(false, TbMsgType.PROVISION_FAILURE),
ASSIGNED_TO_EDGE(false, TbMsgType.ENTITY_ASSIGNED_TO_EDGE), // log edge name
UNASSIGNED_FROM_EDGE(false, TbMsgType.ENTITY_UNASSIGNED_FROM_EDGE),
ADDED_COMMENT(false, TbMsgType.COMMENT_CREATED),
UPDATED_COMMENT(false, TbMsgType.COMMENT_UPDATED),
DELETED_COMMENT(false, null),
SMS_SENT(false, null);
@Getter
private final boolean isRead;
ActionType(boolean isRead) {
private final TbMsgType ruleEngineMsgType;
ActionType(boolean isRead, TbMsgType ruleEngineMsgType) {
this.isRead = isRead;
this.ruleEngineMsgType = ruleEngineMsgType;
}
public Optional<TbMsgType> getRuleEngineMsgType() {
return Optional.ofNullable(ruleEngineMsgType);
}
}

95
common/data/src/main/java/org/thingsboard/server/common/data/msg/TbMsgType.java

@ -0,0 +1,95 @@
/**
* Copyright © 2016-2023 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.msg;
import lombok.Getter;
import java.util.Arrays;
import java.util.EnumSet;
import java.util.List;
import java.util.Objects;
import java.util.stream.Collectors;
public enum TbMsgType {
POST_ATTRIBUTES_REQUEST("Post attributes"),
POST_TELEMETRY_REQUEST("Post telemetry"),
TO_SERVER_RPC_REQUEST("RPC Request from Device"),
ACTIVITY_EVENT("Activity Event"),
INACTIVITY_EVENT("Inactivity Event"),
CONNECT_EVENT("Connect Event"),
DISCONNECT_EVENT("Disconnect Event"),
ENTITY_CREATED("Entity Created"),
ENTITY_UPDATED("Entity Updated"),
ENTITY_DELETED("Entity Deleted"),
ENTITY_ASSIGNED("Entity Assigned"),
ENTITY_UNASSIGNED("Entity Unassigned"),
ATTRIBUTES_UPDATED("Attributes Updated"),
ATTRIBUTES_DELETED("Attributes Deleted"),
ALARM(null),
ALARM_ACK("Alarm Acknowledged"),
ALARM_CLEAR("Alarm Cleared"),
ALARM_DELETE("Alarm Deleted"),
ALARM_ASSIGNED("Alarm Assigned"),
ALARM_UNASSIGNED("Alarm Unassigned"),
COMMENT_CREATED("Comment Created"),
COMMENT_UPDATED("Comment Updated"),
RPC_CALL_FROM_SERVER_TO_DEVICE("RPC Request to Device"),
ENTITY_ASSIGNED_FROM_TENANT("Entity Assigned From Tenant"),
ENTITY_ASSIGNED_TO_TENANT("Entity Assigned To Tenant"),
ENTITY_ASSIGNED_TO_EDGE(null),
ENTITY_UNASSIGNED_FROM_EDGE(null),
TIMESERIES_UPDATED("Timeseries Updated"),
TIMESERIES_DELETED("Timeseries Deleted"),
RPC_QUEUED("RPC Queued"),
RPC_SENT("RPC Sent"),
RPC_DELIVERED("RPC Delivered"),
RPC_SUCCESSFUL("RPC Successful"),
RPC_TIMEOUT("RPC Timeout"),
RPC_EXPIRED("RPC Expired"),
RPC_FAILED("RPC Failed"),
RPC_DELETED("RPC Deleted"),
RELATION_ADD_OR_UPDATE("Relation Added or Updated"),
RELATION_DELETED("Relation Deleted"),
RELATIONS_DELETED("All Relations Deleted"),
PROVISION_SUCCESS(null),
PROVISION_FAILURE(null);
public static final String OTHER = "Other";
public static final List<String> NODE_CONNECTIONS = EnumSet.allOf(TbMsgType.class).stream()
.map(TbMsgType::getNodeConnection).filter(Objects::nonNull).collect(Collectors.toUnmodifiableList());
@Getter
private final String nodeConnection;
TbMsgType(String nodeConnection) {
this.nodeConnection = nodeConnection;
}
public static String getNodeConnection(String msgType) {
if (msgType == null) {
return OTHER;
} else {
return Arrays.stream(TbMsgType.values())
.filter(type -> type.name().equals(msgType))
.findFirst()
.map(TbMsgType::getNodeConnection)
.orElse(OTHER);
}
}
}

3
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/EmptyNodeConfiguration.java

@ -24,7 +24,6 @@ public class EmptyNodeConfiguration implements NodeConfiguration<EmptyNodeConfig
@Override
public EmptyNodeConfiguration defaultConfiguration() {
EmptyNodeConfiguration configuration = new EmptyNodeConfiguration();
return configuration;
return new EmptyNodeConfiguration();
}
}

9
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbRelationTypes.java → rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNodeConnectionType.java

@ -18,9 +18,12 @@ package org.thingsboard.rule.engine.api;
/**
* Created by ashvayka on 19.01.18.
*/
public final class TbRelationTypes {
public final class TbNodeConnectionType {
public static String SUCCESS = "Success";
public static String FAILURE = "Failure";
public static final String SUCCESS = "Success";
public static final String FAILURE = "Failure";
public static final String TRUE = "True";
public static final String FALSE = "False";
}

15
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java

@ -24,6 +24,7 @@ import org.thingsboard.rule.engine.api.ScriptEngine;
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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.script.ScriptLanguage;
@ -31,6 +32,10 @@ import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_CLEAR;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_UPDATED;
@Slf4j
@ -55,13 +60,13 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
withCallback(processAlarm(ctx, msg),
alarmResult -> {
if (alarmResult.alarm == null) {
ctx.tellNext(msg, "False");
ctx.tellNext(msg, TbNodeConnectionType.FALSE);
} else if (alarmResult.isCreated) {
tellNext(ctx, msg, alarmResult, DataConstants.ENTITY_CREATED, "Created");
tellNext(ctx, msg, alarmResult, ENTITY_CREATED.name(), "Created");
} else if (alarmResult.isUpdated) {
tellNext(ctx, msg, alarmResult, DataConstants.ENTITY_UPDATED, "Updated");
tellNext(ctx, msg, alarmResult, ENTITY_UPDATED.name(), "Updated");
} else if (alarmResult.isCleared) {
tellNext(ctx, msg, alarmResult, DataConstants.ALARM_CLEAR, "Cleared");
tellNext(ctx, msg, alarmResult, ALARM_CLEAR.name(), "Cleared");
} else {
ctx.tellSuccess(msg);
}
@ -96,7 +101,7 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
} else if (alarmResult.isCleared) {
metaData.putValue(DataConstants.IS_CLEARED_ALARM, Boolean.TRUE.toString());
}
return ctx.transformMsg(originalMsg, "ALARM", originalMsg.getOriginator(), metaData, data);
return ctx.transformMsg(originalMsg, ALARM.name(), originalMsg.getOriginator(), metaData, data);
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java

@ -57,8 +57,8 @@ import java.util.Optional;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.FAILURE;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
@Slf4j
public abstract class TbAbstractRelationActionNode<C extends TbAbstractRelationActionNodeConfiguration> implements TbNode {

22
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCopyAttributesToEntityViewNode.java

@ -36,7 +36,6 @@ import org.thingsboard.server.common.data.objects.AttributesEntityView;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.util.CollectionsUtil;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import javax.annotation.Nullable;
@ -45,7 +44,12 @@ import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
@Slf4j
@RuleNode(
@ -71,14 +75,14 @@ public class TbCopyAttributesToEntityViewNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (DataConstants.ATTRIBUTES_UPDATED.equals(msg.getType()) ||
DataConstants.ATTRIBUTES_DELETED.equals(msg.getType()) ||
DataConstants.ACTIVITY_EVENT.equals(msg.getType()) ||
DataConstants.INACTIVITY_EVENT.equals(msg.getType()) ||
SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msg.getType())) {
if (ATTRIBUTES_UPDATED.name().equals(msg.getType()) ||
ATTRIBUTES_DELETED.name().equals(msg.getType()) ||
ACTIVITY_EVENT.name().equals(msg.getType()) ||
INACTIVITY_EVENT.name().equals(msg.getType()) ||
POST_ATTRIBUTES_REQUEST.name().equals(msg.getType())) {
if (!msg.getMetaData().getData().isEmpty()) {
long now = System.currentTimeMillis();
String scope = msg.getType().equals(SessionMsgType.POST_ATTRIBUTES_REQUEST.name()) ?
String scope = msg.getType().equals(POST_ATTRIBUTES_REQUEST.name()) ?
DataConstants.CLIENT_SCOPE : msg.getMetaData().getValue(DataConstants.SCOPE);
ListenableFuture<List<EntityView>> entityViewsFuture =
@ -90,7 +94,7 @@ public class TbCopyAttributesToEntityViewNode implements TbNode {
long startTime = entityView.getStartTimeMs();
long endTime = entityView.getEndTimeMs();
if ((endTime != 0 && endTime > now && startTime < now) || (endTime == 0 && startTime < now)) {
if (DataConstants.ATTRIBUTES_DELETED.equals(msg.getType())) {
if (ATTRIBUTES_DELETED.name().equals(msg.getType())) {
List<String> attributes = new ArrayList<>();
for (JsonElement element : new JsonParser().parse(msg.getData()).getAsJsonObject().get("attributes").getAsJsonArray()) {
if (element.isJsonPrimitive()) {

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbMsgCountNode.java

@ -33,7 +33,7 @@ import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
@Slf4j
@RuleNode(

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java

@ -41,7 +41,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
@Slf4j
@RuleNode(

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java

@ -24,7 +24,7 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -196,7 +196,7 @@ public class TbMsgDeduplicationNode implements TbNode {
private void enqueueForTellNextWithRetry(TbContext ctx, TbMsg msg, int retryAttempt) {
if (config.getMaxRetries() > retryAttempt) {
ctx.enqueueForTellNext(msg, TbRelationTypes.SUCCESS,
ctx.enqueueForTellNext(msg, TbNodeConnectionType.SUCCESS,
() -> {
log.trace("[{}][{}][{}] Successfully enqueue deduplication result message!", ctx.getSelfId(), msg.getOriginator(), retryAttempt);
},

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java

@ -32,7 +32,7 @@ import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
@Slf4j
@RuleNode(

52
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java

@ -30,13 +30,23 @@ import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.CONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST;
import static org.thingsboard.server.common.data.msg.TbMsgType.TIMESERIES_UPDATED;
@Slf4j
public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfiguration, S, U> implements TbNode {
@ -73,7 +83,7 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
protected S buildEvent(TbMsg msg, TbContext ctx) {
String msgType = msg.getType();
if (DataConstants.ALARM.equals(msgType)) {
if (ALARM.name().equals(msgType)) {
EdgeEventActionType actionType = getAlarmActionType(msg);
return buildEvent(ctx.getTenantId(), actionType, getUUIDFromMsgData(msg), getAlarmEventType(), null);
} else {
@ -150,19 +160,19 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
protected EdgeEventActionType getEdgeEventActionTypeByMsgType(String msgType, Map<String, String> metadata) {
EdgeEventActionType actionType;
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)
|| DataConstants.TIMESERIES_UPDATED.equals(msgType)) {
if (POST_TELEMETRY_REQUEST.name().equals(msgType)
|| TIMESERIES_UPDATED.name().equals(msgType)) {
actionType = EdgeEventActionType.TIMESERIES_UPDATED;
} else if (DataConstants.ATTRIBUTES_UPDATED.equals(msgType)) {
} else if (ATTRIBUTES_UPDATED.name().equals(msgType)) {
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
} else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
} else if (POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
actionType = EdgeEventActionType.POST_ATTRIBUTES;
} else if (DataConstants.ATTRIBUTES_DELETED.equals(msgType)) {
} else if (ATTRIBUTES_DELETED.name().equals(msgType)) {
actionType = EdgeEventActionType.ATTRIBUTES_DELETED;
} else if (DataConstants.CONNECT_EVENT.equals(msgType)
|| DataConstants.DISCONNECT_EVENT.equals(msgType)
|| DataConstants.ACTIVITY_EVENT.equals(msgType)
|| DataConstants.INACTIVITY_EVENT.equals(msgType)) {
} else if (CONNECT_EVENT.name().equals(msgType)
|| DISCONNECT_EVENT.name().equals(msgType)
|| ACTIVITY_EVENT.name().equals(msgType)
|| INACTIVITY_EVENT.name().equals(msgType)) {
String scope = metadata.get(SCOPE);
if ( StringUtils.isEmpty(scope)) {
actionType = EdgeEventActionType.TIMESERIES_UPDATED;
@ -177,16 +187,16 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
}
protected boolean isSupportedMsgType(String msgType) {
return SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)
|| SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)
|| DataConstants.ATTRIBUTES_UPDATED.equals(msgType)
|| DataConstants.ATTRIBUTES_DELETED.equals(msgType)
|| DataConstants.TIMESERIES_UPDATED.equals(msgType)
|| DataConstants.ALARM.equals(msgType)
|| DataConstants.CONNECT_EVENT.equals(msgType)
|| DataConstants.DISCONNECT_EVENT.equals(msgType)
|| DataConstants.ACTIVITY_EVENT.equals(msgType)
|| DataConstants.INACTIVITY_EVENT.equals(msgType);
return POST_TELEMETRY_REQUEST.name().equals(msgType)
|| POST_ATTRIBUTES_REQUEST.name().equals(msgType)
|| ATTRIBUTES_UPDATED.name().equals(msgType)
|| ATTRIBUTES_DELETED.name().equals(msgType)
|| TIMESERIES_UPDATED.name().equals(msgType)
|| ALARM.name().equals(msgType)
|| CONNECT_EVENT.name().equals(msgType)
|| DISCONNECT_EVENT.name().equals(msgType)
|| ACTIVITY_EVENT.name().equals(msgType)
|| INACTIVITY_EVENT.name().equals(msgType);
}
protected boolean isSupportedOriginator(EntityType entityType) {

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbAssetTypeSwitchNode.java

@ -34,7 +34,8 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
relationTypes = {},
configClazz = EmptyNodeConfiguration.class,
nodeDescription = "Route incoming messages based on the name of the asset profile",
nodeDetails = "Route incoming messages based on the name of the asset profile. The asset profile name is case-sensitive",
nodeDetails = "Route incoming messages based on the name of the asset profile. The asset profile name is case-sensitive<br><br>" +
"Output connection types: Profile name of message originator or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbNodeEmptyConfig")
public class TbAssetTypeSwitchNode extends TbAbstractTypeSwitchNode {

37
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNode.java

@ -18,7 +18,6 @@ package org.thingsboard.rule.engine.filter;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode;
@ -26,26 +25,27 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import javax.annotation.Nullable;
import java.io.IOException;
@Slf4j
@RuleNode(
type = ComponentType.FILTER,
name = "check alarm status",
name = "alarm status filter",
configClazz = TbCheckAlarmStatusNodeConfig.class,
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
nodeDescription = "Checks alarm status.",
nodeDetails = "Checks the alarm status to match one of the specified statuses.",
nodeDetails = "Checks the alarm status to match one of the specified statuses.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeCheckAlarmStatusConfig")
public class TbCheckAlarmStatusNode implements TbNode {
private TbCheckAlarmStatusNodeConfig config;
@Override
@ -60,33 +60,24 @@ public class TbCheckAlarmStatusNode implements TbNode {
ListenableFuture<Alarm> latest = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), alarm.getId());
Futures.addCallback(latest, new FutureCallback<Alarm>() {
Futures.addCallback(latest, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable Alarm result) {
if (result != null) {
boolean isPresent = false;
for (AlarmStatus alarmStatus : config.getAlarmStatusList()) {
if (result.getStatus() == alarmStatus) {
isPresent = true;
break;
}
}
if (isPresent) {
ctx.tellNext(msg, "True");
} else {
ctx.tellNext(msg, "False");
}
} else {
if (result == null) {
ctx.tellFailure(msg, new TbNodeException("No such alarm found."));
return;
}
boolean isPresent = config.getAlarmStatusList().stream()
.anyMatch(alarmStatus -> result.getStatus() == alarmStatus);
ctx.tellNext(msg, isPresent ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
}
@Override
public void onFailure(Throwable t) {
ctx.tellFailure(msg, t);
}
}, MoreExecutors.directExecutor());
} catch (IllegalArgumentException e) {
}, ctx.getDbCallbackExecutor());
} catch (Exception e) {
log.error("Failed to parse alarm: [{}]", msg.getData());
throw new TbNodeException(e);
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNodeConfig.java

@ -24,12 +24,14 @@ import java.util.List;
@Data
public class TbCheckAlarmStatusNodeConfig implements NodeConfiguration<TbCheckAlarmStatusNodeConfig> {
private List<AlarmStatus> alarmStatusList;
@Override
public TbCheckAlarmStatusNodeConfig defaultConfiguration() {
TbCheckAlarmStatusNodeConfig config = new TbCheckAlarmStatusNodeConfig();
var config = new TbCheckAlarmStatusNodeConfig();
config.setAlarmStatusList(Arrays.asList(AlarmStatus.ACTIVE_ACK, AlarmStatus.ACTIVE_UNACK));
return config;
}
}

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckMessageNode.java

@ -22,6 +22,7 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -33,12 +34,12 @@ import java.util.Map;
@RuleNode(
type = ComponentType.FILTER,
name = "check fields presence",
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
configClazz = TbCheckMessageNodeConfiguration.class,
nodeDescription = "Checks the presence of the specified fields in the message and/or metadata.",
nodeDetails = "Checks the presence of the specified fields in the message and/or metadata. " +
"By default, the rule node checks that all specified fields need to be present. " +
"Uncheck the 'Check that all specified fields are present' if the presence of at least one field is sufficient.",
nodeDetails = "By default, the rule node checks that all specified fields are present. " +
"Uncheck the 'Check that all selected fields are present' if the presence of at least one field is sufficient.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeCheckMessageConfig")
public class TbCheckMessageNode implements TbNode {
@ -60,9 +61,9 @@ public class TbCheckMessageNode implements TbNode {
public void onMsg(TbContext ctx, TbMsg msg) {
try {
if (config.isCheckAllKeys()) {
ctx.tellNext(msg, allKeysData(msg) && allKeysMetadata(msg) ? "True" : "False");
ctx.tellNext(msg, allKeysData(msg) && allKeysMetadata(msg) ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
} else {
ctx.tellNext(msg, atLeastOneData(msg) || atLeastOneMetadata(msg) ? "True" : "False");
ctx.tellNext(msg, atLeastOneData(msg) || atLeastOneMetadata(msg) ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
}
} catch (Exception e) {
ctx.tellFailure(msg, e);

29
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNode.java

@ -24,6 +24,7 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -43,12 +44,14 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback;
@Slf4j
@RuleNode(
type = ComponentType.FILTER,
name = "check relation",
name = "check relation presence",
configClazz = TbCheckRelationNodeConfiguration.class,
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
nodeDescription = "Checks the presence of the relation between the originator of the message and other entities.",
nodeDetails = "If 'check relation to specific entity' is selected, one must specify a related entity. " +
"Otherwise, the rule node checks the presence of a relation to any entity that matches the direction and relation type criteria.",
nodeDetails = "If 'check relation to specific entity' is selected, you should specify a related entity. " +
"Otherwise, the rule node checks the presence of a relation to any entity. " +
"In both cases, relation lookup is based on configured direction and type.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeCheckRelationConfig")
public class TbCheckRelationNode implements TbNode {
@ -67,13 +70,11 @@ public class TbCheckRelationNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
ListenableFuture<Boolean> checkRelationFuture;
if (config.isCheckForSingleEntity()) {
checkRelationFuture = processSingle(ctx, msg);
} else {
checkRelationFuture = processList(ctx, msg);
}
withCallback(checkRelationFuture, filterResult -> ctx.tellNext(msg, filterResult ? "True" : "False"), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
ListenableFuture<Boolean> checkRelationFuture = config.isCheckForSingleEntity() ?
processSingle(ctx, msg) : processList(ctx, msg);
withCallback(checkRelationFuture,
filterResult -> ctx.tellNext(msg, filterResult ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
private ListenableFuture<Boolean> processSingle(TbContext ctx, TbMsg msg) {
@ -100,11 +101,7 @@ public class TbCheckRelationNode implements TbNode {
}
private ListenableFuture<Boolean> isEmptyList(List<EntityRelation> entityRelations) {
if (entityRelations.isEmpty()) {
return Futures.immediateFuture(false);
} else {
return Futures.immediateFuture(true);
}
return entityRelations.isEmpty() ? Futures.immediateFuture(false) : Futures.immediateFuture(true);
}
}

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbDeviceTypeSwitchNode.java

@ -34,7 +34,8 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
relationTypes = {"default"},
configClazz = EmptyNodeConfiguration.class,
nodeDescription = "Route incoming messages based on the name of the device profile",
nodeDetails = "Route incoming messages based on the name of the device profile. The device profile name is case-sensitive",
nodeDetails = "Route incoming messages based on the name of the device profile. The device profile name is case-sensitive<br><br>" +
"Output connection types: Profile name of message originator or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbNodeEmptyConfig")
public class TbDeviceTypeSwitchNode extends TbAbstractTypeSwitchNode {

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java

@ -22,6 +22,7 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.script.ScriptLanguage;
@ -32,7 +33,8 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback;
@Slf4j
@RuleNode(
type = ComponentType.FILTER,
name = "script", relationTypes = {"True", "False"},
name = "script",
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
configClazz = TbJsFilterNodeConfiguration.class,
nodeDescription = "Filter incoming messages using TBEL or JS script",
nodeDetails = "Evaluates boolean function using incoming message. " +
@ -40,7 +42,8 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback;
"Script function should return boolean value and accepts three parameters: <br/>" +
"Message payload can be accessed via <code>msg</code> property. For example <code>msg.temperature < 10;</code><br/>" +
"Message metadata can be accessed via <code>metadata</code> property. For example <code>metadata.customerName === 'John';</code><br/>" +
"Message type can be accessed via <code>msgType</code> property.",
"Message type can be accessed via <code>msgType</code> property.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeScriptConfig"
)
@ -62,7 +65,7 @@ public class TbJsFilterNode implements TbNode {
withCallback(scriptEngine.executeFilterAsync(msg),
filterResult -> {
ctx.logJsEvalResponse();
ctx.tellNext(msg, filterResult ? "True" : "False");
ctx.tellNext(msg, filterResult ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
},
t -> {
ctx.tellFailure(msg, t);

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java

@ -44,7 +44,8 @@ import java.util.Set;
"If Array is empty - message not routed to next Node. " +
"Message payload can be accessed via <code>msg</code> property. For example <code>msg.temperature < 10;</code><br/>" +
"Message metadata can be accessed via <code>metadata</code> property. For example <code>metadata.customerName === 'John';</code><br/>" +
"Message type can be accessed via <code>msgType</code> property.",
"Message type can be accessed via <code>msgType</code> property.<br><br>" +
"Output connection types: Custom connection(s) defined by switch node or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeSwitchConfig")
public class TbJsSwitchNode implements TbNode {

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java

@ -21,6 +21,7 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -33,9 +34,10 @@ import org.thingsboard.server.common.msg.TbMsg;
type = ComponentType.FILTER,
name = "message type",
configClazz = TbMsgTypeFilterNodeConfiguration.class,
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
nodeDescription = "Filter incoming messages by Message Type",
nodeDetails = "If incoming MessageType is expected - send Message via <b>True</b> chain, otherwise <b>False</b> chain is used.",
nodeDetails = "If incoming message type is expected - send Message via <b>True</b> chain, otherwise <b>False</b> chain is used.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeMessageTypeConfig")
public class TbMsgTypeFilterNode implements TbNode {
@ -49,7 +51,7 @@ public class TbMsgTypeFilterNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
ctx.tellNext(msg, config.getMessageTypes().contains(msg.getType()) ? "True" : "False");
ctx.tellNext(msg, config.getMessageTypes().contains(msg.getType()) ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
}
}

86
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java

@ -19,24 +19,20 @@ import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.msg.TbMsgType;
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.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType;
@Slf4j
@RuleNode(
type = ComponentType.FILTER,
name = "message type switch",
configClazz = EmptyNodeConfiguration.class,
relationTypes = {"Post attributes", "Post telemetry", "RPC Request from Device", "RPC Request to Device", "RPC Queued", "RPC Sent", "RPC Delivered", "RPC Successful", "RPC Timeout", "RPC Expired", "RPC Failed", "RPC Deleted",
"Activity Event", "Inactivity Event", "Connect Event", "Disconnect Event", "Entity Created", "Entity Updated", "Entity Deleted", "Entity Assigned",
"Entity Unassigned", "Attributes Updated", "Attributes Deleted", "Alarm Acknowledged", "Alarm Cleared", "Alarm Assigned", "Alarm Unassigned", "Comment Created", "Comment Updated", "Other", "Entity Assigned From Tenant", "Entity Assigned To Tenant",
"Relation Added or Updated", "Relation Deleted", "All Relations Deleted", "Timeseries Updated", "Timeseries Deleted"},
relationTypes = {}, // should always be empty. We add the relation types for this node in AnnotationComponentDiscoveryService.
nodeDescription = "Route incoming messages by Message Type",
nodeDetails = "Sends messages with message types <b>\"Post attributes\", \"Post telemetry\", \"RPC Request\"</b> etc. via corresponding chain, otherwise <b>Other</b> chain is used.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
@ -52,83 +48,7 @@ public class TbMsgTypeSwitchNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
String relationType;
if (msg.getType().equals(SessionMsgType.POST_ATTRIBUTES_REQUEST.name())) {
relationType = "Post attributes";
} else if (msg.getType().equals(SessionMsgType.POST_TELEMETRY_REQUEST.name())) {
relationType = "Post telemetry";
} else if (msg.getType().equals(SessionMsgType.TO_SERVER_RPC_REQUEST.name())) {
relationType = "RPC Request from Device";
} else if (msg.getType().equals(DataConstants.ACTIVITY_EVENT)) {
relationType = "Activity Event";
} else if (msg.getType().equals(DataConstants.INACTIVITY_EVENT)) {
relationType = "Inactivity Event";
} else if (msg.getType().equals(DataConstants.CONNECT_EVENT)) {
relationType = "Connect Event";
} else if (msg.getType().equals(DataConstants.DISCONNECT_EVENT)) {
relationType = "Disconnect Event";
} else if (msg.getType().equals(DataConstants.ENTITY_CREATED)) {
relationType = "Entity Created";
} else if (msg.getType().equals(DataConstants.ENTITY_UPDATED)) {
relationType = "Entity Updated";
} else if (msg.getType().equals(DataConstants.ENTITY_DELETED)) {
relationType = "Entity Deleted";
} else if (msg.getType().equals(DataConstants.ENTITY_ASSIGNED)) {
relationType = "Entity Assigned";
} else if (msg.getType().equals(DataConstants.ENTITY_UNASSIGNED)) {
relationType = "Entity Unassigned";
} else if (msg.getType().equals(DataConstants.ATTRIBUTES_UPDATED)) {
relationType = "Attributes Updated";
} else if (msg.getType().equals(DataConstants.ATTRIBUTES_DELETED)) {
relationType = "Attributes Deleted";
} else if (msg.getType().equals(DataConstants.ALARM_ACK)) {
relationType = "Alarm Acknowledged";
} else if (msg.getType().equals(DataConstants.ALARM_CLEAR)) {
relationType = "Alarm Cleared";
} else if (msg.getType().equals(DataConstants.ALARM_ASSIGNED)) {
relationType = "Alarm Assigned";
} else if (msg.getType().equals(DataConstants.ALARM_UNASSIGNED)) {
relationType = "Alarm Unassigned";
} else if (msg.getType().equals(DataConstants.COMMENT_CREATED)) {
relationType = "Comment Created";
} else if (msg.getType().equals(DataConstants.COMMENT_UPDATED)) {
relationType = "Comment Updated";
} else if (msg.getType().equals(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE)) {
relationType = "RPC Request to Device";
} else if (msg.getType().equals(DataConstants.ENTITY_ASSIGNED_FROM_TENANT)) {
relationType = "Entity Assigned From Tenant";
} else if (msg.getType().equals(DataConstants.ENTITY_ASSIGNED_TO_TENANT)) {
relationType = "Entity Assigned To Tenant";
} else if (msg.getType().equals(DataConstants.TIMESERIES_UPDATED)) {
relationType = "Timeseries Updated";
} else if (msg.getType().equals(DataConstants.TIMESERIES_DELETED)) {
relationType = "Timeseries Deleted";
} else if (msg.getType().equals(DataConstants.RPC_QUEUED)) {
relationType = "RPC Queued";
} else if (msg.getType().equals(DataConstants.RPC_SENT)) {
relationType = "RPC Sent";
} else if (msg.getType().equals(DataConstants.RPC_DELIVERED)) {
relationType = "RPC Delivered";
} else if (msg.getType().equals(DataConstants.RPC_SUCCESSFUL)) {
relationType = "RPC Successful";
} else if (msg.getType().equals(DataConstants.RPC_TIMEOUT)) {
relationType = "RPC Timeout";
} else if (msg.getType().equals(DataConstants.RPC_EXPIRED)) {
relationType = "RPC Expired";
} else if (msg.getType().equals(DataConstants.RPC_FAILED)) {
relationType = "RPC Failed";
} else if (msg.getType().equals(DataConstants.RPC_DELETED)) {
relationType = "RPC Deleted";
} else if (msg.getType().equals(DataConstants.RELATION_ADD_OR_UPDATE)) {
relationType = "Relation Added or Updated";
} else if (msg.getType().equals(DataConstants.RELATION_DELETED)) {
relationType = "Relation Deleted";
} else if (msg.getType().equals(DataConstants.RELATIONS_DELETED)) {
relationType = "All Relations Deleted";
} else {
relationType = "Other";
}
ctx.tellNext(msg, relationType);
ctx.tellNext(msg, TbMsgType.getNodeConnection(msg.getType()));
}
}

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeFilterNode.java

@ -21,6 +21,7 @@ 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.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -31,9 +32,10 @@ import org.thingsboard.server.common.msg.TbMsg;
type = ComponentType.FILTER,
name = "entity type",
configClazz = TbOriginatorTypeFilterNodeConfiguration.class,
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
nodeDescription = "Filter incoming messages by the type of message originator entity",
nodeDetails = "Checks that the entity type of the incoming message originator matches one of the values specified in the filter.",
nodeDetails = "Checks that the entity type of the incoming message originator matches one of the values specified in the filter.<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeOriginatorTypeConfig")
public class TbOriginatorTypeFilterNode implements TbNode {
@ -48,7 +50,7 @@ public class TbOriginatorTypeFilterNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
EntityType originatorType = msg.getOriginator().getEntityType();
ctx.tellNext(msg, config.getOriginatorTypes().contains(originatorType) ? "True" : "False");
ctx.tellNext(msg, config.getOriginatorTypes().contains(originatorType) ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
}
}

50
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java

@ -19,8 +19,6 @@ import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -29,55 +27,17 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
type = ComponentType.FILTER,
name = "entity type switch",
configClazz = EmptyNodeConfiguration.class,
relationTypes = {"Device", "Asset", "Alarm", "Entity View", "Tenant", "Customer", "User", "Dashboard", "Rule chain", "Rule node", "Edge"},
relationTypes = {}, // should always be empty. We add the relation types for this node in AnnotationComponentDiscoveryService.
nodeDescription = "Route incoming messages by Message Originator Type",
nodeDetails = "Routes messages to chain according to the entity type ('Device', 'Asset', etc.).",
nodeDetails = "Routes messages to chain according to the entity type ('Device', 'Asset', etc.).<br><br>" +
"Output connection types: entityType of the message originator or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbNodeEmptyConfig")
public class TbOriginatorTypeSwitchNode extends TbAbstractTypeSwitchNode {
@Override
protected String getRelationType(TbContext ctx, EntityId originator) throws TbNodeException {
String relationType;
EntityType originatorType = originator.getEntityType();
switch (originatorType) {
case TENANT:
relationType = "Tenant";
break;
case CUSTOMER:
relationType = "Customer";
break;
case USER:
relationType = "User";
break;
case DASHBOARD:
relationType = "Dashboard";
break;
case ASSET:
relationType = "Asset";
break;
case DEVICE:
relationType = "Device";
break;
case ENTITY_VIEW:
relationType = "Entity View";
break;
case EDGE:
relationType = "Edge";
break;
case RULE_CHAIN:
relationType = "Rule chain";
break;
case RULE_NODE:
relationType = "Rule node";
break;
case ALARM:
relationType = "Alarm";
break;
default:
throw new TbNodeException("Unsupported originator type: " + originatorType);
}
return relationType;
protected String getRelationType(TbContext ctx, EntityId originator) {
return originator.getEntityType().getNormalName();
}
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbCheckpointNode.java

@ -21,7 +21,7 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -48,7 +48,7 @@ public class TbCheckpointNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
ctx.enqueueForTellNext(msg, queueName, TbRelationTypes.SUCCESS, () -> ctx.ack(msg), error -> ctx.tellFailure(msg, error));
ctx.enqueueForTellNext(msg, queueName, TbNodeConnectionType.SUCCESS, () -> ctx.ack(msg), error -> ctx.tellFailure(msg, error));
}
}

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/geo/TbGpsGeofencingFilterNode.java

@ -19,6 +19,7 @@ import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -30,7 +31,7 @@ import org.thingsboard.server.common.msg.TbMsg;
type = ComponentType.FILTER,
name = "gps geofencing filter",
configClazz = TbGpsGeofencingFilterNodeConfiguration.class,
relationTypes = {"True", "False"},
relationTypes = {TbNodeConnectionType.TRUE, TbNodeConnectionType.FALSE},
nodeDescription = "Filter incoming messages by GPS based geofencing",
nodeDetails = "Extracts latitude and longitude parameters from the incoming message and checks them according to configured perimeter. </br>" +
"Configuration:</br></br>" +
@ -57,14 +58,15 @@ import org.thingsboard.server.common.msg.TbMsg;
"</br></br>" +
"{\"latitude\": 48.198618758582384, \"longitude\": 24.65322245153503, \"radius\": 100.0, \"radiusUnit\": \"METER\" }" +
"</br></br>" +
"Available radius units: METER, KILOMETER, FOOT, MILE, NAUTICAL_MILE;",
"Available radius units: METER, KILOMETER, FOOT, MILE, NAUTICAL_MILE;<br><br>" +
"Output connection types: <code>True</code>, <code>False</code>, <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeGpsGeofencingConfig")
public class TbGpsGeofencingFilterNode extends AbstractGeofencingNode<TbGpsGeofencingFilterNodeConfiguration> {
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
ctx.tellNext(msg, checkMatches(msg) ? "True" : "False");
ctx.tellNext(msg, checkMatches(msg) ? TbNodeConnectionType.TRUE : TbNodeConnectionType.FALSE);
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java

@ -31,7 +31,7 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.exception.ThingsboardKafkaClientError;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -165,7 +165,7 @@ public class TbKafkaNode implements TbNode {
private void processRecord(TbContext ctx, TbMsg msg, RecordMetadata metadata, Exception e) {
if (e == null) {
TbMsg next = processResponse(ctx, msg, metadata);
ctx.tellNext(next, TbRelationTypes.SUCCESS);
ctx.tellNext(next, TbNodeConnectionType.SUCCESS);
} else {
TbMsg next = processException(ctx, msg, e);
ctx.tellFailure(next, e);

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java

@ -27,7 +27,6 @@ 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.StringUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -35,7 +34,7 @@ import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.TbNodeConnectionType.SUCCESS;
import static org.thingsboard.rule.engine.mail.TbSendEmailNode.SEND_EMAIL_TYPE;
@Slf4j

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java

@ -49,7 +49,7 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
protected void checkIfDetailsListIsNotEmptyOrElseThrow(List<ContactBasedEntityDetails> detailsList) throws TbNodeException {
if (detailsList == null || detailsList.isEmpty()) {
throw new TbNodeException("No entity details selected!");
throw new TbNodeException("At least one entity detail should be selected!");
}
}

31
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java

@ -40,7 +40,6 @@ import org.thingsboard.server.common.data.query.EntityKey;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.rule.RuleNodeState;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.dao.sql.query.EntityKeyMapping;
@ -55,6 +54,18 @@ import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_ACK;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_CLEAR;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_DELETE;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_ASSIGNED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_UNASSIGNED;
import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST;
@Slf4j
class DeviceState {
@ -136,24 +147,24 @@ class DeviceState {
latestValues = fetchLatestValues(ctx, deviceId);
}
boolean stateChanged = false;
if (msg.getType().equals(SessionMsgType.POST_TELEMETRY_REQUEST.name())) {
if (msg.getType().equals(POST_TELEMETRY_REQUEST.name())) {
stateChanged = processTelemetry(ctx, msg);
} else if (msg.getType().equals(SessionMsgType.POST_ATTRIBUTES_REQUEST.name())) {
} else if (msg.getType().equals(POST_ATTRIBUTES_REQUEST.name())) {
stateChanged = processAttributesUpdateRequest(ctx, msg);
} else if (msg.getType().equals(DataConstants.ACTIVITY_EVENT) || msg.getType().equals(DataConstants.INACTIVITY_EVENT)) {
} else if (msg.getType().equals(ACTIVITY_EVENT.name()) || msg.getType().equals(INACTIVITY_EVENT.name())) {
stateChanged = processDeviceActivityEvent(ctx, msg);
} else if (msg.getType().equals(DataConstants.ATTRIBUTES_UPDATED)) {
} else if (msg.getType().equals(ATTRIBUTES_UPDATED.name())) {
stateChanged = processAttributesUpdateNotification(ctx, msg);
} else if (msg.getType().equals(DataConstants.ATTRIBUTES_DELETED)) {
} else if (msg.getType().equals(ATTRIBUTES_DELETED.name())) {
stateChanged = processAttributesDeleteNotification(ctx, msg);
} else if (msg.getType().equals(DataConstants.ALARM_CLEAR)) {
} else if (msg.getType().equals(ALARM_CLEAR.name())) {
stateChanged = processAlarmClearNotification(ctx, msg);
} else if (msg.getType().equals(DataConstants.ALARM_ACK)) {
} else if (msg.getType().equals(ALARM_ACK.name())) {
processAlarmAckNotification(ctx, msg);
} else if (msg.getType().equals(DataConstants.ALARM_DELETE)) {
} else if (msg.getType().equals(ALARM_DELETE.name())) {
processAlarmDeleteNotification(ctx, msg);
} else {
if (msg.getType().equals(DataConstants.ENTITY_ASSIGNED) || msg.getType().equals(DataConstants.ENTITY_UNASSIGNED)) {
if (msg.getType().equals(ENTITY_ASSIGNED.name()) || msg.getType().equals(ENTITY_UNASSIGNED.name())) {
dynamicPredicateValueCtx.resetCustomer();
}
ctx.tellSuccess(msg);

8
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java

@ -26,7 +26,6 @@ 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.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
@ -46,6 +45,9 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_UPDATED;
@Slf4j
@RuleNode(
type = ComponentType.ACTION,
@ -123,10 +125,10 @@ public class TbDeviceProfileNode implements TbNode {
} else {
if (EntityType.DEVICE.equals(originatorType)) {
DeviceId deviceId = new DeviceId(msg.getOriginator().getId());
if (msg.getType().equals(DataConstants.ENTITY_UPDATED)) {
if (msg.getType().equals(ENTITY_UPDATED.name())) {
invalidateDeviceProfileCache(deviceId, msg.getData());
ctx.tellSuccess(msg);
} else if (msg.getType().equals(DataConstants.ENTITY_DELETED)) {
} else if (msg.getType().equals(ENTITY_DELETED.name())) {
removeDeviceState(deviceId);
ctx.tellSuccess(msg);
} else {

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java

@ -43,7 +43,7 @@ import org.springframework.web.util.UriComponentsBuilder;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.credentials.BasicCredentials;
import org.thingsboard.rule.engine.credentials.ClientCredentials;
@ -212,7 +212,7 @@ public class TbHttpClient {
ctx.tellSuccess(next);
} else {
TbMsg next = processFailureResponse(ctx, msg, responseEntity);
ctx.tellNext(next, TbRelationTypes.FAILURE);
ctx.tellNext(next, TbNodeConnectionType.FAILURE);
}
}
});

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java

@ -27,12 +27,13 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -76,7 +77,7 @@ public class TbSendRPCRequestNode implements TbNode {
ctx.tellFailure(msg, new RuntimeException("Params are not present in the message!"));
} else {
int requestId = json.has("requestId") ? json.get("requestId").getAsInt() : random.nextInt();
boolean restApiCall = msg.getType().equals(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE);
boolean restApiCall = msg.getType().equals(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE.name());
tmp = msg.getMetaData().getValue("oneway");
boolean oneway = !StringUtils.isEmpty(tmp) && Boolean.parseBoolean(tmp);
@ -117,7 +118,7 @@ public class TbSendRPCRequestNode implements TbNode {
ctx.getRpcService().sendRpcRequestToDevice(request, ruleEngineDeviceRpcResponse -> {
if (ruleEngineDeviceRpcResponse.getError().isEmpty()) {
TbMsg next = ctx.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getCustomerId(), msg.getMetaData(), ruleEngineDeviceRpcResponse.getResponse().orElse("{}"));
ctx.enqueueForTellNext(next, TbRelationTypes.SUCCESS);
ctx.enqueueForTellNext(next, TbNodeConnectionType.SUCCESS);
} else {
TbMsg next = ctx.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getCustomerId(), msg.getMetaData(), wrap("error", ruleEngineDeviceRpcResponse.getError().get().name()));
ctx.enqueueForTellFailure(next, ruleEngineDeviceRpcResponse.getError().get().name());

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java

@ -22,7 +22,7 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.TbMsgCallback;
@ -75,7 +75,7 @@ public abstract class TbAbstractTransformNode<C> implements TbNode {
ctx.tellFailure(msg, e);
}
});
msgs.forEach(newMsg -> ctx.enqueueForTellNext(newMsg, TbRelationTypes.SUCCESS, wrapper::onSuccess, wrapper::onFailure));
msgs.forEach(newMsg -> ctx.enqueueForTellNext(newMsg, TbNodeConnectionType.SUCCESS, wrapper::onSuccess, wrapper::onFailure));
}
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java

@ -25,7 +25,7 @@ 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.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -79,7 +79,7 @@ public class TbSplitArrayMsgNode implements TbNode {
});
data.forEach(msgNode -> {
TbMsg outMsg = TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode));
ctx.enqueueForTellNext(outMsg, TbRelationTypes.SUCCESS, wrapper::onSuccess, wrapper::onFailure);
ctx.enqueueForTellNext(outMsg, TbNodeConnectionType.SUCCESS, wrapper::onSuccess, wrapper::onFailure);
});
}
} else {

13
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbCreateRelationNodeTest.java

@ -28,8 +28,8 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
@ -54,6 +54,7 @@ import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.msg.TbMsgType.ENTITY_CREATED;
@RunWith(MockitoJUnitRunner.class)
public class TbCreateRelationNodeTest {
@ -111,7 +112,7 @@ public class TbCreateRelationNodeTest {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("name", "AssetName");
metaData.putValue("type", "AssetType");
msg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
msg = TbMsg.newMsg(ENTITY_CREATED.name(), deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(ctx.getRelationService().checkRelationAsync(any(), eq(assetId), eq(deviceId), eq(RELATION_TYPE_CONTAINS), eq(RelationTypeGroup.COMMON)))
.thenReturn(Futures.immediateFuture(false));
@ -119,7 +120,7 @@ public class TbCreateRelationNodeTest {
.thenReturn(Futures.immediateFuture(true));
node.onMsg(ctx, msg);
verify(ctx).tellNext(msg, TbRelationTypes.SUCCESS);
verify(ctx).tellNext(msg, TbNodeConnectionType.SUCCESS);
}
@Test
@ -138,7 +139,7 @@ public class TbCreateRelationNodeTest {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("name", "AssetName");
metaData.putValue("type", "AssetType");
msg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
msg = TbMsg.newMsg(ENTITY_CREATED.name(), deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
EntityRelation relation = new EntityRelation();
when(ctx.getRelationService().findByToAndTypeAsync(any(), eq(msg.getOriginator()), eq(RELATION_TYPE_CONTAINS), eq(RelationTypeGroup.COMMON)))
@ -150,7 +151,7 @@ public class TbCreateRelationNodeTest {
.thenReturn(Futures.immediateFuture(true));
node.onMsg(ctx, msg);
verify(ctx).tellNext(msg, TbRelationTypes.SUCCESS);
verify(ctx).tellNext(msg, TbNodeConnectionType.SUCCESS);
}
@Test
@ -169,7 +170,7 @@ public class TbCreateRelationNodeTest {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("name", "AssetName");
metaData.putValue("type", "AssetType");
msg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
msg = TbMsg.newMsg(ENTITY_CREATED.name(), deviceId, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(ctx.getRelationService().checkRelationAsync(any(), eq(assetId), eq(deviceId), eq(RELATION_TYPE_CONTAINS), eq(RelationTypeGroup.COMMON)))
.thenReturn(Futures.immediateFuture(false));

18
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNodeTest.java

@ -49,10 +49,18 @@ import java.util.UUID;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.verify;
import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.common.data.msg.TbMsgType.CONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT;
@RunWith(MockitoJUnitRunner.class)
public class TbMsgPushToEdgeNodeTest {
private static final List<String> MISC_EVENTS = List.of(CONNECT_EVENT.name(), DISCONNECT_EVENT.name(),
ACTIVITY_EVENT.name(), INACTIVITY_EVENT.name());
TbMsgPushToEdgeNode node;
private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
@ -102,7 +110,7 @@ public class TbMsgPushToEdgeNodeTest {
PageData<EdgeId> edgePageData = new PageData<>(List.of(edgeId), 1, 1, false);
Mockito.when(edgeService.findRelatedEdgeIdsByEntityId(tenantId, userId, new PageLink(TbMsgPushToEdgeNode.DEFAULT_PAGE_SIZE))).thenReturn(edgePageData);
TbMsg msg = TbMsg.newMsg(DataConstants.ATTRIBUTES_UPDATED, userId, new TbMsgMetaData(),
TbMsg msg = TbMsg.newMsg(ATTRIBUTES_UPDATED.name(), userId, new TbMsgMetaData(),
TbMsgDataType.JSON, "{}", null, null);
node.onMsg(ctx, msg);
@ -112,9 +120,7 @@ public class TbMsgPushToEdgeNodeTest {
@Test
public void testMiscEventsProcessedAsAttributesUpdated() {
List<String> miscEvents = List.of(DataConstants.CONNECT_EVENT, DataConstants.DISCONNECT_EVENT,
DataConstants.ACTIVITY_EVENT, DataConstants.INACTIVITY_EVENT);
for (String event : miscEvents) {
for (String event : MISC_EVENTS) {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue(DataConstants.SCOPE, DataConstants.SERVER_SCOPE);
testEvent(event, metaData, EdgeEventActionType.ATTRIBUTES_UPDATED, "kv");
@ -123,9 +129,7 @@ public class TbMsgPushToEdgeNodeTest {
@Test
public void testMiscEventsProcessedAsTimeseriesUpdated() {
List<String> miscEvents = List.of(DataConstants.CONNECT_EVENT, DataConstants.DISCONNECT_EVENT,
DataConstants.ACTIVITY_EVENT, DataConstants.INACTIVITY_EVENT);
for (String event : miscEvents) {
for (String event : MISC_EVENTS) {
testEvent(event, new TbMsgMetaData(), EdgeEventActionType.TIMESERIES_UPDATED, "data");
}
}

14
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsFilterNodeTest.java

@ -27,6 +27,8 @@ import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.script.ScriptLanguage;
@ -55,21 +57,21 @@ public class TbJsFilterNodeTest {
private RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased());
@Test
public void falseEvaluationDoNotSendMsg() throws TbNodeException, ScriptException {
public void falseEvaluationDoNotSendMsg() throws TbNodeException {
initWithScript();
TbMsg msg = TbMsg.newMsg("USER", null, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
TbMsg msg = TbMsg.newMsg(EntityType.USER.name(), null, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(scriptEngine.executeFilterAsync(msg)).thenReturn(Futures.immediateFuture(false));
node.onMsg(ctx, msg);
verify(ctx).getDbCallbackExecutor();
verify(ctx).tellNext(msg, "False");
verify(ctx).tellNext(msg, TbNodeConnectionType.FALSE);
}
@Test
public void exceptionInJsThrowsException() throws TbNodeException {
initWithScript();
TbMsgMetaData metaData = new TbMsgMetaData();
TbMsg msg = TbMsg.newMsg("USER", null, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
TbMsg msg = TbMsg.newMsg(EntityType.USER.name(), null, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(scriptEngine.executeFilterAsync(msg)).thenReturn(Futures.immediateFailedFuture(new ScriptException("error")));
@ -81,12 +83,12 @@ public class TbJsFilterNodeTest {
public void metadataConditionCanBeTrue() throws TbNodeException {
initWithScript();
TbMsgMetaData metaData = new TbMsgMetaData();
TbMsg msg = TbMsg.newMsg("USER", null, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
TbMsg msg = TbMsg.newMsg(EntityType.USER.name(), null, metaData, TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(scriptEngine.executeFilterAsync(msg)).thenReturn(Futures.immediateFuture(true));
node.onMsg(ctx, msg);
verify(ctx).getDbCallbackExecutor();
verify(ctx).tellNext(msg, "True");
verify(ctx).tellNext(msg, TbNodeConnectionType.TRUE);
}
private void initWithScript() throws TbNodeException {

15
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/DeviceStateTest.java

@ -22,7 +22,6 @@ import org.mockito.ArgumentCaptor;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult;
@ -64,6 +63,10 @@ import static org.mockito.Mockito.never;
import static org.mockito.Mockito.reset;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_CLEAR;
import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM_DELETE;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST;
public class DeviceStateTest {
@ -115,11 +118,11 @@ public class DeviceStateTest {
verify(ctx).enqueueForTellNext(resultMsgCaptor.capture(), eq("Alarm Created"));
Alarm alarm = JacksonUtil.fromString(resultMsgCaptor.getValue().getData(), Alarm.class);
deviceState.process(ctx, TbMsg.newMsg(DataConstants.ALARM_CLEAR, deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm)));
deviceState.process(ctx, TbMsg.newMsg(ALARM_CLEAR.name(), deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm)));
reset(ctx);
String deletedAttributes = "{ \"attributes\": [ \"other\" ] }";
deviceState.process(ctx, TbMsg.newMsg(DataConstants.ATTRIBUTES_DELETED, deviceId, new TbMsgMetaData(), deletedAttributes));
deviceState.process(ctx, TbMsg.newMsg(ATTRIBUTES_DELETED.name(), deviceId, new TbMsgMetaData(), deletedAttributes));
verify(ctx, never()).enqueueForTellNext(any(), anyString());
}
@ -129,7 +132,7 @@ public class DeviceStateTest {
DeviceId deviceId = new DeviceId(UUID.randomUUID());
DeviceState deviceState = createDeviceState(deviceId, alarmConfig);
TbMsg attributeUpdateMsg = TbMsg.newMsg(SessionMsgType.POST_ATTRIBUTES_REQUEST.name(),
TbMsg attributeUpdateMsg = TbMsg.newMsg(POST_ATTRIBUTES_REQUEST.name(),
deviceId, new TbMsgMetaData(), "{ \"enabled\": false }");
deviceState.process(ctx, attributeUpdateMsg);
@ -137,9 +140,9 @@ public class DeviceStateTest {
verify(ctx).enqueueForTellNext(resultMsgCaptor.capture(), eq("Alarm Created"));
Alarm alarm = JacksonUtil.fromString(resultMsgCaptor.getValue().getData(), Alarm.class);
deviceState.process(ctx, TbMsg.newMsg(DataConstants.ALARM_CLEAR, deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm)));
deviceState.process(ctx, TbMsg.newMsg(ALARM_CLEAR.name(), deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm)));
TbMsg alarmDeleteNotification = TbMsg.newMsg(DataConstants.ALARM_DELETE, deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm));
TbMsg alarmDeleteNotification = TbMsg.newMsg(ALARM_DELETE.name(), deviceId, new TbMsgMetaData(), JacksonUtil.toString(alarm));
assertDoesNotThrow(() -> {
deviceState.process(ctx, alarmDeleteNotification);
});

12
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java

@ -30,7 +30,7 @@ import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.TbNodeConnectionType;
import org.thingsboard.rule.engine.deduplication.DeduplicationStrategy;
import org.thingsboard.rule.engine.deduplication.TbMsgDeduplicationNode;
import org.thingsboard.rule.engine.deduplication.TbMsgDeduplicationNodeConfiguration;
@ -173,7 +173,7 @@ public class TbMsgDeduplicationNodeTest {
verify(ctx, times(msgCount)).ack(any());
verify(ctx, times(1)).tellFailure(eq(msgToReject), any());
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbNodeConnectionType.SUCCESS), successCaptor.capture(), failureCaptor.capture());
TbMsg firstMsg = inputMsgs.get(0);
TbMsg actualMsg = newMsgCaptor.getValue();
@ -221,7 +221,7 @@ public class TbMsgDeduplicationNodeTest {
verify(ctx, times(msgCount)).ack(any());
verify(ctx, times(1)).tellFailure(eq(msgToReject), any());
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbNodeConnectionType.SUCCESS), successCaptor.capture(), failureCaptor.capture());
TbMsg actualMsg = newMsgCaptor.getValue();
// msg ids should be different because we create new msg before enqueueForTellNext
@ -263,7 +263,7 @@ public class TbMsgDeduplicationNodeTest {
verify(ctx, times(msgCount)).ack(any());
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture());
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbNodeConnectionType.SUCCESS), successCaptor.capture(), failureCaptor.capture());
Assertions.assertEquals(1, newMsgCaptor.getAllValues().size());
TbMsg outMessage = newMsgCaptor.getAllValues().get(0);
@ -309,7 +309,7 @@ public class TbMsgDeduplicationNodeTest {
verify(ctx, times(msgCount)).ack(any());
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any());
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture());
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbNodeConnectionType.SUCCESS), successCaptor.capture(), failureCaptor.capture());
List<TbMsg> resultMsgs = newMsgCaptor.getAllValues();
Assertions.assertEquals(2, resultMsgs.size());
@ -363,7 +363,7 @@ public class TbMsgDeduplicationNodeTest {
verify(ctx, times(msgCount)).ack(any());
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any());
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture());
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbNodeConnectionType.SUCCESS), successCaptor.capture(), failureCaptor.capture());
List<TbMsg> resultMsgs = newMsgCaptor.getAllValues();
Assertions.assertEquals(2, resultMsgs.size());

Loading…
Cancel
Save