Browse Source

added checkMsgType util method to TbMsg & resolved other review comments

pull/8786/head
ShvaykaD 3 years ago
parent
commit
2d4fbd6833
  1. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  2. 14
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  3. 44
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  4. 6
      common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java
  5. 22
      common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java
  6. 54
      common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionMsgType.java
  7. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java
  8. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCopyAttributesToEntityViewNode.java
  9. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbMsgCountNode.java
  10. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java
  11. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java
  12. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java
  13. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java
  14. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/delay/TbMsgDelayNode.java
  15. 51
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java
  16. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbAssetTypeSwitchNode.java
  17. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNode.java
  18. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbDeviceTypeSwitchNode.java
  19. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java
  20. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java
  21. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java
  22. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java
  23. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java
  24. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  25. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  26. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java
  27. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java
  28. 31
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java
  29. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java
  30. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java
  31. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbHttpClient.java
  32. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
  33. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  34. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java

4
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -223,9 +223,9 @@ public class DefaultTbClusterService implements TbClusterService {
if (isRuleChainTransform && isQueueTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId, targetQueueName);
} else if (isRuleChainTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId);
tbMsg = TbMsg.transformMsgRuleChainId(tbMsg, targetRuleChainId);
} else if (isQueueTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetQueueName);
tbMsg = TbMsg.transformMsgQueueName(tbMsg, targetQueueName);
}
return tbMsg;
}

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

@ -36,7 +36,6 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceIdInfo;
import org.thingsboard.server.common.data.EntityType;
@ -103,6 +102,9 @@ import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
/**
* Created by ashvayka on 01.05.18.
*/
@ -575,7 +577,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
ListenableFuture<List<TsKvEntry>> tsData = tsService.findLatest(TenantId.SYS_TENANT_ID, device.getId(), PERSISTENT_ATTRIBUTES);
future = Futures.transform(tsData, extractDeviceStateData(device), deviceStateExecutor);
} else {
ListenableFuture<List<AttributeKvEntry>> attrData = attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), DataConstants.SERVER_SCOPE, PERSISTENT_ATTRIBUTES);
ListenableFuture<List<AttributeKvEntry>> attrData = attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), SERVER_SCOPE, PERSISTENT_ATTRIBUTES);
future = Futures.transform(attrData, extractDeviceStateData(device), deviceStateExecutor);
}
return transformInactivityTimeout(future);
@ -586,7 +588,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
if (!persistToTelemetry || deviceStateData.getState().getInactivityTimeout() != defaultInactivityTimeoutMs) {
return future; //fail fast
}
var attributesFuture = attributesService.find(TenantId.SYS_TENANT_ID, deviceStateData.getDeviceId(), DataConstants.SERVER_SCOPE, INACTIVITY_TIMEOUT);
var attributesFuture = attributesService.find(TenantId.SYS_TENANT_ID, deviceStateData.getDeviceId(), SERVER_SCOPE, INACTIVITY_TIMEOUT);
return Futures.transform(attributesFuture, attributes -> {
attributes.flatMap(KvEntry::getLongValue).ifPresent((inactivityTimeout) -> {
if (inactivityTimeout > 0) {
@ -779,7 +781,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
}
TbMsgMetaData md = stateData.getMetaData().copy();
if (!persistToTelemetry) {
md.putValue(DataConstants.SCOPE, DataConstants.SERVER_SCOPE);
md.putValue(SCOPE, SERVER_SCOPE);
}
TbMsg tbMsg = TbMsg.newMsg(msgType, stateData.getDeviceId(), stateData.getCustomerId(), md, TbMsgDataType.JSON, data);
clusterService.pushMsgToRuleEngine(stateData.getTenantId(), stateData.getDeviceId(), tbMsg, null);
@ -795,7 +797,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value));
} else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
}
}
@ -806,7 +808,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value));
} else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
}
}

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

@ -54,9 +54,53 @@ 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";

6
common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java

@ -38,15 +38,15 @@ public class StringUtils {
}
public static boolean isBlank(String source) {
return isEmpty(source) || source.trim().isEmpty();
return source == null || source.isEmpty() || source.trim().isEmpty();
}
public static boolean isNotEmpty(String source) {
return !isEmpty(source);
return source != null && !source.isEmpty();
}
public static boolean isNotBlank(String source) {
return !isBlank(source);
return source != null && !source.isEmpty() && !source.trim().isEmpty();
}
public static String notBlankOrDefault(String src, String def) {

22
common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java

@ -279,7 +279,7 @@ public final class TbMsg implements Serializable {
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, TbMsgMetaData metadata) {
public static TbMsg transformMsgMetadata(TbMsg tbMsg, TbMsgMetaData metadata) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, metadata.copy(), tbMsg.dataType,
tbMsg.data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.getCallback());
}
@ -289,17 +289,17 @@ public final class TbMsg implements Serializable {
data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, CustomerId customerId) {
public static TbMsg transformMsgCustomerId(TbMsg tbMsg, CustomerId customerId) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.ruleChainId, tbMsg.ruleNodeId, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, RuleChainId ruleChainId) {
public static TbMsg transformMsgRuleChainId(TbMsg tbMsg, RuleChainId ruleChainId) {
return new TbMsg(tbMsg.queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, ruleChainId, null, tbMsg.ctx.copy(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, String queueName) {
public static TbMsg transformMsgQueueName(TbMsg tbMsg, String queueName) {
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.customerId, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.getRuleChainId(), null, tbMsg.ctx.copy(), tbMsg.getCallback());
}
@ -467,4 +467,18 @@ public final class TbMsg implements Serializable {
}
return ts;
}
public boolean checkType(TbMsgType tbMsgType) {
return tbMsgType != null && tbMsgType.name().equals(this.type);
}
public boolean checkTypeOneOf(TbMsgType... types) {
for (TbMsgType type : types) {
if (checkType(type)) {
return true;
}
}
return false;
}
}

54
common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionMsgType.java

@ -0,0 +1,54 @@
/**
* 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.msg.session;
/**
* @deprecated This enum is deprecated and will be removed in a future version.
* Note: This enum was originally part of the public API but is now specific to CoAP transport only.
* Please use {@link org.thingsboard.server.transport.coap.CoapSessionMsgType} instead.
*/
@Deprecated(since="3.5.2", forRemoval = true)
public enum SessionMsgType {
GET_ATTRIBUTES_REQUEST(true), POST_ATTRIBUTES_REQUEST(true), GET_ATTRIBUTES_RESPONSE,
SUBSCRIBE_ATTRIBUTES_REQUEST, UNSUBSCRIBE_ATTRIBUTES_REQUEST, ATTRIBUTES_UPDATE_NOTIFICATION,
POST_TELEMETRY_REQUEST(true), STATUS_CODE_RESPONSE,
SUBSCRIBE_RPC_COMMANDS_REQUEST, UNSUBSCRIBE_RPC_COMMANDS_REQUEST,
TO_DEVICE_RPC_REQUEST, TO_DEVICE_RPC_RESPONSE, TO_DEVICE_RPC_RESPONSE_ACK,
TO_SERVER_RPC_REQUEST(true), TO_SERVER_RPC_RESPONSE,
RULE_ENGINE_ERROR,
SESSION_OPEN, SESSION_CLOSE,
CLAIM_REQUEST();
private final boolean requiresRulesProcessing;
SessionMsgType() {
this(false);
}
SessionMsgType(boolean requiresRulesProcessing) {
this.requiresRulesProcessing = requiresRulesProcessing;
}
public boolean requiresRulesProcessing() {
return requiresRulesProcessing;
}
}

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

@ -24,10 +24,10 @@ 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.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.data.script.ScriptLanguage;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
@ -79,7 +79,7 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
if (previousDetails != null) {
TbMsgMetaData metaData = msg.getMetaData().copy();
metaData.putValue(PREV_ALARM_DETAILS, JacksonUtil.toString(previousDetails));
dummyMsg = TbMsg.transformMsg(msg, metaData);
dummyMsg = TbMsg.transformMsgMetadata(msg, metaData);
}
return scriptEngine.executeJsonAsync(dummyMsg);
} catch (Exception e) {

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

@ -75,14 +75,11 @@ public class TbCopyAttributesToEntityViewNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
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.checkTypeOneOf(ATTRIBUTES_UPDATED, ATTRIBUTES_DELETED,
ACTIVITY_EVENT, INACTIVITY_EVENT, POST_ATTRIBUTES_REQUEST)) {
if (!msg.getMetaData().getData().isEmpty()) {
long now = System.currentTimeMillis();
String scope = msg.getType().equals(POST_ATTRIBUTES_REQUEST.name()) ?
String scope = msg.checkType(POST_ATTRIBUTES_REQUEST) ?
DataConstants.CLIENT_SCOPE : msg.getMetaData().getValue(DataConstants.SCOPE);
ListenableFuture<List<EntityView>> entityViewsFuture =
@ -94,7 +91,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 (ATTRIBUTES_DELETED.name().equals(msg.getType())) {
if (msg.checkType(ATTRIBUTES_DELETED)) {
List<String> attributes = new ArrayList<>();
for (JsonElement element : JsonParser.parseString(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

@ -65,7 +65,7 @@ public class TbMsgCountNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (msg.getType().equals(TbMsgType.MSG_COUNT_SELF_MSG.name()) && msg.getId().equals(nextTickId)) {
if (msg.checkType(TbMsgType.MSG_COUNT_SELF_MSG) && msg.getId().equals(nextTickId)) {
JsonObject telemetryJson = new JsonObject();
telemetryJson.addProperty(this.telemetryPrefix + "_" + ctx.getServiceId(), messagesProcessed.longValue());

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java

@ -105,13 +105,13 @@ public class TbSnsNode extends TbAbstractExternalNode {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(MESSAGE_ID, result.getMessageId());
metaData.putValue(REQUEST_ID, result.getSdkResponseMetadata().getRequestId());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private TbMsg processException(TbContext ctx, TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
@Override

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java

@ -134,13 +134,13 @@ public class TbSqsNode extends TbAbstractExternalNode {
if (!StringUtils.isEmpty(result.getSequenceNumber())) {
metaData.putValue(SEQUENCE_NUMBER, result.getSequenceNumber());
}
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
@Override

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

@ -107,7 +107,7 @@ public class TbMsgGeneratorNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
log.trace("onMsg, config {}, msg {}", config, msg);
if (initialized.get() && msg.getType().equals(TbMsgType.GENERATOR_NODE_SELF_MSG.name()) && msg.getId().equals(nextTickId)) {
if (initialized.get() && msg.checkType(TbMsgType.GENERATOR_NODE_SELF_MSG) && msg.getId().equals(nextTickId)) {
TbStopWatch sw = TbStopWatch.create();
withCallback(generate(ctx, msg),
m -> {
@ -146,7 +146,7 @@ public class TbMsgGeneratorNode implements TbNode {
private ListenableFuture<TbMsg> generate(TbContext ctx, TbMsg msg) {
log.trace("generate, config {}", config);
if (prevMsg == null) {
prevMsg = ctx.newMsg(config.getQueueName(), "", originatorId, msg.getCustomerId(), TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
prevMsg = ctx.newMsg(config.getQueueName(), TbMsg.EMPTY_STRING, originatorId, msg.getCustomerId(), TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
}
if (initialized.get()) {
ctx.logJsEvalRequest();

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

@ -24,10 +24,10 @@ 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.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg;
@ -80,7 +80,7 @@ public class TbMsgDeduplicationNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException {
if (TbMsgType.DEDUPLICATION_TIMEOUT_SELF_MSG.name().equals(msg.getType())) {
if (msg.checkType(TbMsgType.DEDUPLICATION_TIMEOUT_SELF_MSG)) {
processDeduplication(ctx, msg.getOriginator());
} else {
processOnRegularMsg(ctx, msg);

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

@ -61,7 +61,7 @@ public class TbMsgDelayNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (msg.getType().equals(TbMsgType.DELAY_TIMEOUT_SELF_MSG.name())) {
if (msg.checkType(TbMsgType.DELAY_TIMEOUT_SELF_MSG)) {
TbMsg pendingMsg = pendingMsgs.remove(UUID.fromString(msg.getData()));
if (pendingMsg != null) {
ctx.enqueueForTellNext(

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

@ -67,7 +67,7 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
return;
}
if (isSupportedOriginator(msg.getOriginator().getEntityType())) {
if (isSupportedMsgType(msg.getType())) {
if (isSupportedMsgType(msg)) {
processMsg(ctx, msg);
} else {
String errMsg = String.format("Unsupported msg type %s", msg.getType());
@ -82,13 +82,12 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
}
protected S buildEvent(TbMsg msg, TbContext ctx) {
String msgType = msg.getType();
if (ALARM.name().equals(msgType)) {
if (msg.checkType(ALARM)) {
EdgeEventActionType actionType = getAlarmActionType(msg);
return buildEvent(ctx.getTenantId(), actionType, getUUIDFromMsgData(msg), getAlarmEventType(), null);
} else {
Map<String, String> metadata = msg.getMetaData().getData();
EdgeEventActionType actionType = getEdgeEventActionTypeByMsgType(msgType, metadata);
EdgeEventActionType actionType = getEdgeEventActionTypeByMsgType(msg);
Map<String, Object> entityBody = new HashMap<>();
JsonNode dataJson = JacksonUtil.toJsonNode(msg.getData());
switch (actionType) {
@ -158,45 +157,31 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
return scope;
}
protected EdgeEventActionType getEdgeEventActionTypeByMsgType(String msgType, Map<String, String> metadata) {
protected EdgeEventActionType getEdgeEventActionTypeByMsgType(TbMsg msg) {
EdgeEventActionType actionType;
if (POST_TELEMETRY_REQUEST.name().equals(msgType)
|| TIMESERIES_UPDATED.name().equals(msgType)) {
if (msg.checkTypeOneOf(POST_TELEMETRY_REQUEST, TIMESERIES_UPDATED)) {
actionType = EdgeEventActionType.TIMESERIES_UPDATED;
} else if (ATTRIBUTES_UPDATED.name().equals(msgType)) {
} else if (msg.checkType(ATTRIBUTES_UPDATED)) {
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
} else if (POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
} else if (msg.checkType(POST_ATTRIBUTES_REQUEST)) {
actionType = EdgeEventActionType.POST_ATTRIBUTES;
} else if (ATTRIBUTES_DELETED.name().equals(msgType)) {
} else if (msg.checkType(ATTRIBUTES_DELETED)) {
actionType = EdgeEventActionType.ATTRIBUTES_DELETED;
} 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;
} else {
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
}
} else if (msg.checkTypeOneOf(CONNECT_EVENT, DISCONNECT_EVENT, ACTIVITY_EVENT, INACTIVITY_EVENT)) {
String scope = msg.getMetaData().getValue(SCOPE);
actionType = StringUtils.isEmpty(scope) ?
EdgeEventActionType.TIMESERIES_UPDATED : EdgeEventActionType.ATTRIBUTES_UPDATED;
} else {
log.warn("Unsupported msg type [{}]", msgType);
throw new IllegalArgumentException("Unsupported msg type: " + msgType);
String type = msg.getType();
log.warn("Unsupported msg type [{}]", type);
throw new IllegalArgumentException("Unsupported msg type: " + type);
}
return actionType;
}
protected boolean isSupportedMsgType(String 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 isSupportedMsgType(TbMsg msg) {
return msg.checkTypeOneOf(POST_TELEMETRY_REQUEST, POST_ATTRIBUTES_REQUEST, ATTRIBUTES_UPDATED,
ATTRIBUTES_DELETED, TIMESERIES_UPDATED, ALARM, CONNECT_EVENT, DISCONNECT_EVENT, ACTIVITY_EVENT, INACTIVITY_EVENT);
}
protected boolean isSupportedOriginator(EntityType entityType) {

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

@ -35,7 +35,7 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
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.<br><br>" +
"Output connections: <i>Message originator profile name</i> or <code>Failure</code>",
"Output connections: <i>Asset profile name</i> or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbNodeEmptyConfig")
public class TbAssetTypeSwitchNode extends TbAbstractTypeSwitchNode {

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

@ -118,11 +118,11 @@ public class TbCheckRelationNode implements TbVersionedNode {
throw new TbNodeException("property to update: '" + DIRECTION_PROPERTY_NAME + "' doesn't exists in configuration!");
}
String direction = newConfigObjectNode.get(DIRECTION_PROPERTY_NAME).asText();
if ("TO".equals(direction)) {
if (EntitySearchDirection.TO.name().equals(direction)) {
newConfigObjectNode.put(DIRECTION_PROPERTY_NAME, EntitySearchDirection.FROM.name());
return new TbPair<>(true, newConfigObjectNode);
}
if ("FROM".equals(direction)) {
if (EntitySearchDirection.FROM.name().equals(direction)) {
newConfigObjectNode.put(DIRECTION_PROPERTY_NAME, EntitySearchDirection.TO.name());
return new TbPair<>(true, newConfigObjectNode);
}

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

@ -35,7 +35,7 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
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<br><br>" +
"Output connections: <i>Message originator profile name</i> or <code>Failure</code>",
"Output connections: <i>Device profile name</i> or <code>Failure</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbNodeEmptyConfig")
public class TbDeviceTypeSwitchNode extends TbAbstractTypeSwitchNode {

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java

@ -119,13 +119,13 @@ public class TbPubSubNode extends TbAbstractExternalNode {
private TbMsg processPublishResult(TbMsg origMsg, String messageId) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(MESSAGE_ID, messageId);
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private Publisher initPubSubClient() throws IOException {

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

@ -176,13 +176,13 @@ public class TbKafkaNode extends TbAbstractExternalNode {
metaData.putValue(OFFSET, String.valueOf(recordMetadata.offset()));
metaData.putValue(PARTITION, String.valueOf(recordMetadata.partition()));
metaData.putValue(TOPIC, recordMetadata.topic());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private TbMsg processException(TbMsg origMsg, Exception e) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, e.getClass() + ": " + e.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
}

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java

@ -70,7 +70,7 @@ public class TbSendEmailNode extends TbAbstractExternalNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
try {
validateType(msg.getType());
validateType(msg);
TbEmail email = getEmail(msg);
var tbMsg = ackIfNeeded(ctx, msg);
withCallback(ctx.getMailExecutor().executeAsync(() -> {
@ -100,8 +100,9 @@ public class TbSendEmailNode extends TbAbstractExternalNode {
return email;
}
private void validateType(String type) {
if (!TbMsgType.SEND_EMAIL.name().equals(type)) {
private void validateType(TbMsg msg) {
if (!msg.checkType(TbMsgType.SEND_EMAIL)) {
String type = msg.getType();
log.warn("Not expected msg type [{}] for SendEmail Node", type);
throw new IllegalStateException("Not expected msg type " + type + " for SendEmail Node");
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java

@ -248,7 +248,7 @@ public class TbMathNode implements TbNode {
} else {
md.putValue(mathResultKey, Double.toString(toDoubleValue(mathResultDef, result)));
}
return TbMsg.transformMsg(msg, md);
return TbMsg.transformMsgMetadata(msg, md);
}
private double calculateResult(List<TbMathArgumentValue> args) {

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java

@ -25,12 +25,12 @@ import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
@ -75,7 +75,7 @@ public class CalculateDeltaNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (!msg.getType().equals(TbMsgType.POST_TELEMETRY_REQUEST.name())) {
if (!msg.checkType(TbMsgType.POST_TELEMETRY_REQUEST)) {
ctx.tellNext(msg, TbNodeConnectionType.OTHER);
return;
}

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

@ -83,7 +83,7 @@ public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeC
case DATA:
return TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode));
case METADATA:
return TbMsg.transformMsg(msg, msgMetaData);
return TbMsg.transformMsgMetadata(msg, msgMetaData);
default:
log.debug("Unexpected FetchTo value: {}. Allowed values: {}", fetchTo, FetchTo.values());
return msg;

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

@ -101,7 +101,7 @@ public class TbGetTelemetryNode implements TbNode {
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(interval, keys));
DonAsynchron.withCallback(list, data -> {
var metaData = updateMetadata(data, msg, keys);
ctx.tellSuccess(TbMsg.transformMsg(msg, metaData));
ctx.tellSuccess(TbMsg.transformMsgMetadata(msg, metaData));
}, error -> ctx.tellFailure(msg, error), ctx.getDbCallbackExecutor());
} catch (Exception e) {
ctx.tellFailure(msg, e);

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java

@ -95,7 +95,7 @@ public class TbMqttNode extends TbAbstractExternalNode {
private TbMsg processException(TbMsg origMsg, Throwable e) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, e.getClass() + ": " + e.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
@Override

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java

@ -19,7 +19,6 @@ import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
@ -77,7 +76,7 @@ public class TbNotificationNode extends TbAbstractExternalNode {
ctx.getNotificationCenter().processNotificationRequest(ctx.getTenantId(), notificationRequest, stats -> {
TbMsgMetaData metaData = tbMsg.getMetaData().copy();
metaData.putValue("notificationRequestResult", JacksonUtil.toString(stats));
tellSuccess(ctx, TbMsg.transformMsg(tbMsg, metaData));
tellSuccess(ctx, TbMsg.transformMsgMetadata(tbMsg, metaData));
})),
r -> {
},

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

@ -36,7 +36,6 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.query.EntityKey;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.rule.RuleNodeState;
@ -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(TbMsgType.POST_TELEMETRY_REQUEST.name())) {
if (msg.checkType(POST_TELEMETRY_REQUEST)) {
stateChanged = processTelemetry(ctx, msg);
} else if (msg.getType().equals(TbMsgType.POST_ATTRIBUTES_REQUEST.name())) {
} else if (msg.checkType(POST_ATTRIBUTES_REQUEST)) {
stateChanged = processAttributesUpdateRequest(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ACTIVITY_EVENT.name()) || msg.getType().equals(TbMsgType.INACTIVITY_EVENT.name())) {
} else if (msg.checkTypeOneOf(ACTIVITY_EVENT, INACTIVITY_EVENT)) {
stateChanged = processDeviceActivityEvent(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ATTRIBUTES_UPDATED.name())) {
} else if (msg.checkType(ATTRIBUTES_UPDATED)) {
stateChanged = processAttributesUpdateNotification(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ATTRIBUTES_DELETED.name())) {
} else if (msg.checkType(ATTRIBUTES_DELETED)) {
stateChanged = processAttributesDeleteNotification(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ALARM_CLEAR.name())) {
} else if (msg.checkType(ALARM_CLEAR)) {
stateChanged = processAlarmClearNotification(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ALARM_ACK.name())) {
} else if (msg.checkType(ALARM_ACK)) {
processAlarmAckNotification(ctx, msg);
} else if (msg.getType().equals(TbMsgType.ALARM_DELETE.name())) {
} else if (msg.checkType(ALARM_DELETE)) {
processAlarmDeleteNotification(ctx, msg);
} else {
if (msg.getType().equals(TbMsgType.ENTITY_ASSIGNED.name()) || msg.getType().equals(TbMsgType.ENTITY_UNASSIGNED.name())) {
if (msg.checkTypeOneOf(ENTITY_ASSIGNED, ENTITY_UNASSIGNED)) {
dynamicPredicateValueCtx.resetCustomer();
}
ctx.tellSuccess(msg);

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

@ -108,16 +108,16 @@ public class TbDeviceProfileNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException {
EntityType originatorType = msg.getOriginator().getEntityType();
if (msg.getType().equals(TbMsgType.DEVICE_PROFILE_PERIODIC_SELF_MSG.name())) {
if (msg.checkType(TbMsgType.DEVICE_PROFILE_PERIODIC_SELF_MSG)) {
scheduleAlarmHarvesting(ctx, msg);
harvestAlarms(ctx, System.currentTimeMillis());
return;
}
if (msg.getType().equals(TbMsgType.DEVICE_PROFILE_UPDATE_SELF_MSG.name())) {
if (msg.checkType(TbMsgType.DEVICE_PROFILE_UPDATE_SELF_MSG)) {
updateProfile(ctx, new DeviceProfileId(UUID.fromString(msg.getData())));
return;
}
if (msg.getType().equals(TbMsgType.DEVICE_UPDATE_SELF_MSG.name())) {
if (msg.checkType(TbMsgType.DEVICE_UPDATE_SELF_MSG)) {
JsonNode data = JacksonUtil.toJsonNode(msg.getData());
DeviceId deviceId = new DeviceId(UUID.fromString(data.get("deviceId").asText()));
if (data.has("profileId")) {
@ -129,12 +129,12 @@ public class TbDeviceProfileNode implements TbNode {
}
if (EntityType.DEVICE.equals(originatorType)) {
DeviceId deviceId = new DeviceId(msg.getOriginator().getId());
if (msg.getType().equals(TbMsgType.ENTITY_UPDATED.name())) {
if (msg.checkType(TbMsgType.ENTITY_UPDATED)) {
invalidateDeviceProfileCache(deviceId, msg.getData());
ctx.tellSuccess(msg);
return;
}
if (msg.getType().equals(TbMsgType.ENTITY_DELETED.name())) {
if (msg.checkType(TbMsgType.ENTITY_DELETED)) {
removeDeviceState(deviceId);
ctx.tellSuccess(msg);
return;

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java

@ -118,7 +118,7 @@ public class TbRabbitMqNode extends TbAbstractExternalNode {
private TbMsg processException(TbMsg origMsg, Throwable t) {
TbMsgMetaData metaData = origMsg.getMetaData().copy();
metaData.putValue(ERROR, t.getClass() + ": " + t.getMessage());
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
@Override

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

@ -286,7 +286,7 @@ public class TbHttpClient {
metaData.putValue(STATUS_REASON, response.getStatusCode().getReasonPhrase());
metaData.putValue(ERROR_BODY, response.getBody());
headersToMetaData(response.getHeaders(), metaData::putValue);
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private TbMsg processException(TbMsg origMsg, Throwable e) {
@ -298,7 +298,7 @@ public class TbHttpClient {
metaData.putValue(STATUS_CODE, restClientResponseException.getRawStatusCode() + "");
metaData.putValue(ERROR_BODY, restClientResponseException.getResponseBodyAsString());
}
return TbMsg.transformMsg(origMsg, metaData);
return TbMsg.transformMsgMetadata(origMsg, metaData);
}
private HttpHeaders prepareHeaders(TbMsg msg) {

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

@ -27,13 +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.server.common.data.msg.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.msg.TbNodeConnectionType;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -76,7 +76,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(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE.name());
boolean restApiCall = msg.checkType(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE);
tmp = msg.getMetaData().getValue("oneway");
boolean oneway = !StringUtils.isEmpty(tmp) && Boolean.parseBoolean(tmp);

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

@ -65,7 +65,7 @@ public class TbMsgAttributesNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (!msg.getType().equals(POST_ATTRIBUTES_REQUEST.name())) {
if (!msg.checkType(POST_ATTRIBUTES_REQUEST)) {
ctx.tellFailure(msg, new IllegalArgumentException("Unsupported msg type: " + msg.getType()));
return;
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java

@ -82,7 +82,7 @@ public class TbMsgTimeseriesNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (!msg.getType().equals(POST_TELEMETRY_REQUEST.name())) {
if (!msg.checkType(POST_TELEMETRY_REQUEST)) {
ctx.tellFailure(msg, new IllegalArgumentException("Unsupported msg type: " + msg.getType()));
return;
}

Loading…
Cancel
Save