Browse Source

Add fetch to logic, fix and add tests for some rule nodes

pull/8661/head
Dmytro Skarzhynets 4 years ago
parent
commit
f29e1f5fef
  1. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java
  2. 8
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java
  3. 21
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java
  4. 23
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java
  5. 46
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  6. 110
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java
  7. 25
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java
  8. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java
  9. 50
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  10. 109
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java
  11. 36
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java
  12. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java
  13. 15
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  14. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java
  15. 17
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java
  16. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java
  17. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java
  18. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  19. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java
  20. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java
  21. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java
  22. 80
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java
  23. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java
  24. 25
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java
  25. 18
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java
  26. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java
  27. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java
  28. 42
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java
  29. 20
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java
  30. 58
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java
  31. 6
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java
  32. 81
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  33. 376
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java
  34. 83
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java
  35. 85
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java
  36. 231
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java
  37. 7
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java

@ -16,7 +16,5 @@
package org.thingsboard.rule.engine.api;
public interface NodeConfiguration<T extends NodeConfiguration> {
T defaultConfiguration();
}

8
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java

@ -24,13 +24,13 @@ import java.util.concurrent.ExecutionException;
* Created by ashvayka on 19.01.18.
*/
public interface TbNode {
void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException;
void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException;
default void destroy() {}
default void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) {}
default void destroy() {
}
default void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) {
}
}

21
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java

@ -0,0 +1,21 @@
/**
* 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.rule.engine.metadata;
public enum FetchTo {
DATA,
METADATA
}

23
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java

@ -0,0 +1,23 @@
/**
* 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.rule.engine.metadata;
import lombok.Data;
@Data
public abstract class TbAbstractFetchToNodeConfiguration {
private FetchTo fetchTo;
}

46
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java

@ -22,10 +22,8 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.BooleanUtils;
import org.thingsboard.common.util.JacksonUtil;
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;
@ -53,26 +51,19 @@ import static org.thingsboard.server.common.data.DataConstants.LATEST_TS;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeConfiguration, T extends EntityId> implements TbNode {
public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeConfiguration, T extends EntityId> extends TbAbstractNodeWithFetchTo<C> {
private static final String VALUE = "value";
private static final String TS = "ts";
protected C config;
private boolean fetchToData;
private boolean isTellFailureIfAbsent;
private boolean getLatestValueWithTs;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = loadGetAttributesNodeConfig(configuration);
this.fetchToData = config.isFetchToData();
this.getLatestValueWithTs = config.isGetLatestValueWithTs();
this.isTellFailureIfAbsent = BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true);
super.init(ctx, configuration);
getLatestValueWithTs = config.isGetLatestValueWithTs();
isTellFailureIfAbsent = config.isTellFailureIfAbsent();
}
protected abstract C loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException;
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
try {
@ -92,13 +83,9 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
ctx.tellNext(msg, FAILURE);
return;
}
JsonNode msgDataNode;
if (fetchToData) {
msgDataNode = JacksonUtil.toJsonNode(msg.getData());
if (!msgDataNode.isObject()) {
ctx.tellFailure(msg, new IllegalArgumentException("Msg body is not an object!"));
return;
}
ObjectNode msgDataNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataNode = getMsgDataAsObjectNode(msg);
} else {
msgDataNode = null;
}
@ -116,17 +103,22 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
String prefix = getPrefix(keyScope);
kvEntryList.forEach(kvEntry -> {
String key = prefix + kvEntry.getKey();
if (fetchToData) {
JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, key);
} else {
if (FetchTo.DATA.equals(fetchTo)) {
JacksonUtil.addKvEntry(msgDataNode, kvEntry, key);
} else if (FetchTo.METADATA.equals(fetchTo)) {
msgMetaData.putValue(key, kvEntry.getValueAsString());
}
});
});
});
TbMsg outMsg = fetchToData ?
TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode)) :
TbMsg.transformMsg(msg, msgMetaData);
TbMsg outMsg = null;
if (FetchTo.DATA.equals(fetchTo)) {
outMsg = TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode));
} else if (FetchTo.METADATA.equals(fetchTo)) {
outMsg = TbMsg.transformMsg(msg, msgMetaData);
}
if (failuresMap.isEmpty()) {
ctx.tellSuccess(outMsg);
} else {
@ -175,7 +167,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) {
ObjectMapper mapper = fetchToData ? JacksonUtil.OBJECT_MAPPER : JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER;
ObjectMapper mapper = FetchTo.DATA.equals(fetchTo) ? JacksonUtil.OBJECT_MAPPER : JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER;
ObjectNode value = JacksonUtil.newObjectNode(mapper);
value.put(TS, tsKvEntry.getTs());
JacksonUtil.addKvEntry(value, tsKvEntry, VALUE, mapper);

110
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java

@ -0,0 +1,110 @@
/**
* 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.rule.engine.metadata;
import com.fasterxml.jackson.databind.node.ObjectNode;
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.TbContext;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@Slf4j
public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends TbAbstractNodeWithFetchTo<TbGetEntityAttrNodeConfiguration> {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
ObjectNode msgDataAsJsonNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode = getMsgDataAsObjectNode(msg);
} else {
msgDataAsJsonNode = null;
}
ctx.checkTenantEntity(msg.getOriginator());
withCallback(findEntityAsync(ctx, msg.getOriginator()),
entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsJsonNode),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) {
if (entityId == null || entityId.isNullUid()) {
ctx.tellNext(msg, FAILURE);
return;
}
Map<String, String> mappingsMap = new HashMap<>();
config.getAttrMapping().forEach((key, value) -> {
String patternProcessedSourceKey = TbNodeUtils.processPattern(key, msg);
String patternProcessedTargetKey = TbNodeUtils.processPattern(value, msg);
mappingsMap.put(patternProcessedSourceKey, patternProcessedTargetKey);
});
var sourceKeys = List.copyOf(mappingsMap.keySet());
withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys),
data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
var latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrKeys);
return Futures.transform(latest, l ->
l.stream()
.map(i -> (KvEntry) i)
.collect(Collectors.toList()),
MoreExecutors.directExecutor());
}
private ListenableFuture<List<KvEntry>> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
var latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, timeseriesKeys);
return Futures.transform(latest, l ->
l.stream()
.map(i -> (KvEntry) i)
.collect(Collectors.toList()),
MoreExecutors.directExecutor());
}
private void putDataAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> data, Map<String, String> map, ObjectNode msgDataAsJsonNode) {
for (KvEntry entry : data) {
String targetKey = map.get(entry.getKey());
String value = entry.getValueAsString();
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode.put(targetKey, value);
} else if (FetchTo.METADATA.equals(fetchTo)) {
msg.getMetaData().putValue(targetKey, value);
}
}
if (FetchTo.DATA.equals(fetchTo)) {
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataAsJsonNode)));
} else if (FetchTo.METADATA.equals(fetchTo)) {
ctx.tellSuccess(msg);
}
}
}

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

@ -27,11 +27,9 @@ import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
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.util.EntityDetails;
import org.thingsboard.server.common.data.ContactBased;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
@ -41,20 +39,11 @@ import java.util.Map;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
@Slf4j
public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEntityDetailsNodeConfiguration> implements TbNode {
public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEntityDetailsNodeConfiguration> extends TbAbstractNodeWithFetchTo<C> {
private static final Gson gson = new Gson();
private static final JsonParser jsonParser = new JsonParser();
private static final Type TYPE = new TypeToken<Map<String, String>>() {
}.getType();
protected C config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = loadGetEntityDetailsNodeConfiguration(configuration);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
withCallback(getDetails(ctx, msg),
@ -62,22 +51,20 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
protected abstract C loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException;
protected abstract ListenableFuture<TbMsg> getDetails(TbContext ctx, TbMsg msg);
protected abstract ListenableFuture<ContactBased> getContactBasedListenableFuture(TbContext ctx, TbMsg msg);
protected MessageData getDataAsJson(TbMsg msg) {
if (this.config.isAddToMetadata()) {
if (config.getFetchTo() == FetchTo.METADATA) {
return new MessageData(gson.toJsonTree(msg.getMetaData().getData(), TYPE), DataSource.METADATA);
} else {
return new MessageData(jsonParser.parse(msg.getData()), DataSource.DATA);
return new MessageData(JsonParser.parseString(msg.getData()), DataSource.DATA);
}
}
protected ListenableFuture<TbMsg> getTbMsgListenableFuture(TbContext ctx, TbMsg msg, MessageData messageData, String prefix) {
if (this.config.getDetailsList().isEmpty()) {
if (config.getDetailsList().isEmpty()) {
return Futures.immediateFuture(msg);
} else {
ListenableFuture<ContactBased> contactBasedListenableFuture = getContactBasedListenableFuture(ctx, msg);
@ -178,6 +165,4 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
private enum DataSource {
DATA, METADATA
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java

@ -16,16 +16,13 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.util.EntityDetails;
import java.util.List;
@Data
public abstract class TbAbstractGetEntityDetailsNodeConfiguration {
@EqualsAndHashCode(callSuper = true)
public abstract class TbAbstractGetEntityDetailsNodeConfiguration extends TbAbstractFetchToNodeConfiguration {
private List<EntityDetails> detailsList;
private boolean addToMetadata;
}

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

@ -0,0 +1,50 @@
/**
* 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.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import org.thingsboard.common.util.JacksonUtil;
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.msg.TbMsg;
public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeConfiguration> implements TbNode {
protected C config;
protected FetchTo fetchTo;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
config = loadNodeConfiguration(configuration);
if (config.getFetchTo() == null) {
throw new TbNodeException("FetchTo cannot be NULL!");
} else {
fetchTo = config.getFetchTo();
}
}
protected abstract C loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException;
protected ObjectNode getMsgDataAsObjectNode(TbMsg msg) {
JsonNode msgDataNode = JacksonUtil.toJsonNode(msg.getData());
if (!msgDataNode.isObject()) {
throw new IllegalArgumentException("Message body is not an object!");
}
return (ObjectNode) msgDataNode;
}
}

109
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java

@ -1,109 +0,0 @@
/**
* 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.rule.engine.metadata;
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.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.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.msg.TbMsg;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@Slf4j
public abstract class TbEntityGetAttrNode<T extends EntityId> implements TbNode {
private TbGetEntityAttrNodeConfiguration config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
try {
withCallback(findEntityAsync(ctx, msg.getOriginator()),
entityId -> safeGetAttributes(ctx, msg, entityId),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellFailure(msg, th);
}
}
private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId) {
if (entityId == null || entityId.isNullUid()) {
ctx.tellNext(msg, FAILURE);
return;
}
Map<String, String> mappingsMap = new HashMap<>();
config.getAttrMapping().forEach((key, value) -> {
String processPatternKey = TbNodeUtils.processPattern(key, msg);
String processPatternValue = TbNodeUtils.processPattern(value, msg);
mappingsMap.put(processPatternKey, processPatternValue);
});
List<String> keys = List.copyOf(mappingsMap.keySet());
withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, keys) : getAttributesAsync(ctx, entityId, keys),
attributes -> putAttributesAndTell(ctx, msg, attributes, mappingsMap),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrKeys);
return Futures.transform(latest, l ->
l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor());
}
private ListenableFuture<List<KvEntry>> getLatestTelemetry(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, timeseriesKeys);
return Futures.transform(latest, l ->
l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor());
}
private void putAttributesAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> attributes, Map<String, String> map) {
attributes.forEach(r -> {
String attrName = map.get(r.getKey());
msg.getMetaData().putValue(attrName, r.getValueAsString());
});
ctx.tellSuccess(msg);
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
public void setConfig(TbGetEntityAttrNodeConfiguration config) {
this.config = config;
}
}

36
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java

@ -15,24 +15,19 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.concurrent.ExecutionException;
@ -48,39 +43,36 @@ import java.util.concurrent.ExecutionException;
"- send Message via <code>Failure</code> chain, otherwise <code>Success</code> chain is used.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeFetchDeviceCredentialsConfig")
public class TbFetchDeviceCredentialsNode implements TbNode {
public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo<TbFetchDeviceCredentialsNodeConfiguration> {
private static final String CREDENTIALS = "credentials";
private static final String CREDENTIALS_TYPE = "credentialsType";
TbFetchDeviceCredentialsNodeConfiguration config;
boolean fetchToMetadata;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbFetchDeviceCredentialsNodeConfiguration.class);
this.fetchToMetadata = config.isFetchToMetadata();
protected TbFetchDeviceCredentialsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbFetchDeviceCredentialsNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException {
EntityId originator = msg.getOriginator();
var originator = msg.getOriginator();
if (!EntityType.DEVICE.equals(originator.getEntityType())) {
ctx.tellFailure(msg, new RuntimeException("Unsupported originator type: " + originator.getEntityType() + "!"));
return;
}
DeviceId deviceId = new DeviceId(msg.getOriginator().getId());
DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(ctx.getTenantId(), deviceId);
var deviceId = new DeviceId(msg.getOriginator().getId());
var deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(ctx.getTenantId(), deviceId);
if (deviceCredentials == null) {
ctx.tellFailure(msg, new RuntimeException("Failed to get Device Credentials for device: " + deviceId + "!"));
return;
}
TbMsg transformedMsg;
DeviceCredentialsType credentialsType = deviceCredentials.getCredentialsType();
JsonNode credentialsInfo = ctx.getDeviceCredentialsService().toCredentialsInfo(deviceCredentials);
if (fetchToMetadata) {
TbMsgMetaData metaData = msg.getMetaData();
TbMsg transformedMsg = null;
var credentialsType = deviceCredentials.getCredentialsType();
var credentialsInfo = ctx.getDeviceCredentialsService().toCredentialsInfo(deviceCredentials);
if (FetchTo.METADATA.equals(fetchTo)) {
var metaData = msg.getMetaData();
metaData.putValue(CREDENTIALS_TYPE, credentialsType.name());
if (credentialsType.equals(DeviceCredentialsType.ACCESS_TOKEN) || credentialsType.equals(DeviceCredentialsType.X509_CERTIFICATE)) {
metaData.putValue(CREDENTIALS, credentialsInfo.asText());
@ -88,7 +80,7 @@ public class TbFetchDeviceCredentialsNode implements TbNode {
metaData.putValue(CREDENTIALS, JacksonUtil.toString(credentialsInfo));
}
transformedMsg = TbMsg.transformMsg(msg, msg.getType(), originator, metaData, msg.getData());
} else {
} else if (FetchTo.DATA.equals(fetchTo)) {
ObjectNode data = (ObjectNode) JacksonUtil.toJsonNode(msg.getData());
data.put(CREDENTIALS_TYPE, credentialsType.name());
data.set(CREDENTIALS, credentialsInfo);

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java

@ -17,18 +17,17 @@ package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
@Data
@EqualsAndHashCode(callSuper = true)
@JsonIgnoreProperties(ignoreUnknown = true)
public class TbFetchDeviceCredentialsNodeConfiguration implements NodeConfiguration<TbFetchDeviceCredentialsNodeConfiguration> {
private boolean fetchToMetadata;
public class TbFetchDeviceCredentialsNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbFetchDeviceCredentialsNodeConfiguration> {
@Override
public TbFetchDeviceCredentialsNodeConfiguration defaultConfiguration() {
TbFetchDeviceCredentialsNodeConfiguration configuration = new TbFetchDeviceCredentialsNodeConfiguration();
configuration.setFetchToMetadata(true);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

15
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java

@ -32,25 +32,24 @@ import org.thingsboard.server.common.msg.TbMsg;
*/
@Slf4j
@RuleNode(type = ComponentType.ENRICHMENT,
name = "originator attributes",
configClazz = TbGetAttributesNodeConfiguration.class,
nodeDescription = "Enrich the message body or metadata with the originator attributes and/or timeseries data",
nodeDetails = "If Attributes enrichment configured, <b>CLIENT/SHARED/SERVER</b> attributes are added into Message data/metadata " +
name = "originator attributes",
configClazz = TbGetAttributesNodeConfiguration.class,
nodeDescription = "Enrich the message body or metadata with the originator attributes and/or timeseries data",
nodeDetails = "If Attributes enrichment configured, <b>CLIENT/SHARED/SERVER</b> attributes are added into Message data/metadata " +
"with specific prefix: <i>cs/shared/ss</i>. Latest telemetry value added into Message data/metadata without prefix. " +
"To access those attributes in other nodes this template can be used " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.cs_temperature</code> or <code>metadata.shared_limit</code> ",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorAttributesConfig")
public class TbGetAttributesNode extends TbAbstractGetAttributesNode<TbGetAttributesNodeConfiguration, EntityId> {
@Override
protected TbGetAttributesNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetAttributesNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetAttributesNodeConfiguration.class);
}
@Override
protected ListenableFuture<EntityId> findEntityIdAsync(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return Futures.immediateFuture(msg.getOriginator());
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java

@ -16,8 +16,8 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
import java.util.List;
@ -25,8 +25,8 @@ import java.util.List;
* Created by ashvayka on 19.01.18.
*/
@Data
public class TbGetAttributesNodeConfiguration implements NodeConfiguration<TbGetAttributesNodeConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetAttributesNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetAttributesNodeConfiguration> {
private List<String> clientAttributeNames;
private List<String> sharedAttributeNames;
private List<String> serverAttributeNames;
@ -35,7 +35,6 @@ public class TbGetAttributesNodeConfiguration implements NodeConfiguration<TbGet
private boolean tellFailureIfAbsent;
private boolean getLatestValueWithTs;
private boolean fetchToData;
@Override
public TbGetAttributesNodeConfiguration defaultConfiguration() {
@ -46,7 +45,7 @@ public class TbGetAttributesNodeConfiguration implements NodeConfiguration<TbGet
configuration.setLatestTsKeyNames(Collections.emptyList());
configuration.setTellFailureIfAbsent(true);
configuration.setGetLatestValueWithTs(false);
configuration.setFetchToData(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

17
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java

@ -18,6 +18,9 @@ package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.rule.engine.api.RuleNode;
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.util.TbNodeUtils;
import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
@ -25,20 +28,24 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
@RuleNode(
type = ComponentType.ENRICHMENT,
name="customer attributes",
name = "customer attributes",
configClazz = TbGetEntityAttrNodeConfiguration.class,
nodeDescription = "Add Originators Customer Attributes or Latest Telemetry into Message Metadata",
nodeDetails = "Enrich the message metadata with the corresponding customer's latest attributes or telemetry value. " +
nodeDescription = "Add Originators Customer Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "Enrich the Message Metadata/Data with the corresponding customer's latest attributes or telemetry value. " +
"The customer is selected based on the originator of the message: device, asset, etc. " +
"</br>" +
"Useful when you store some parameters on the customer level and would like to use them for message processing.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeCustomerAttributesConfig")
public class TbGetCustomerAttributeNode extends TbEntityGetAttrNode<CustomerId> {
public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode<CustomerId> {
@Override
protected ListenableFuture<CustomerId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator);
}
@Override
protected TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java

@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -48,16 +49,16 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeEntityDetailsConfig")
public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode<TbGetCustomerDetailsNodeConfiguration> {
private static final String CUSTOMER_PREFIX = "customer_";
@Override
protected TbGetCustomerDetailsNodeConfiguration loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetCustomerDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetCustomerDetailsNodeConfiguration.class);
}
@Override
protected ListenableFuture<TbMsg> getDetails(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), CUSTOMER_PREFIX);
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java

@ -16,18 +16,19 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetCustomerDetailsNodeConfiguration extends TbAbstractGetEntityDetailsNodeConfiguration implements NodeConfiguration<TbGetCustomerDetailsNodeConfiguration> {
@Override
public TbGetCustomerDetailsNodeConfiguration defaultConfiguration() {
TbGetCustomerDetailsNodeConfiguration configuration = new TbGetCustomerDetailsNodeConfiguration();
configuration.setDetailsList(Collections.emptyList());
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java

@ -39,15 +39,14 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeDeviceAttributesConfig")
public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode<TbGetDeviceAttrNodeConfiguration, DeviceId> {
@Override
protected TbGetDeviceAttrNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetDeviceAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetDeviceAttrNodeConfiguration.class);
}
@Override
protected ListenableFuture<DeviceId> findEntityIdAsync(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery());
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@ -23,8 +24,8 @@ import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfiguration {
private DeviceRelationsQuery deviceRelationsQuery;
@Override
@ -36,7 +37,7 @@ public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfigu
configuration.setLatestTsKeyNames(Collections.emptyList());
configuration.setTellFailureIfAbsent(true);
configuration.setGetLatestValueWithTs(false);
configuration.setFetchToData(false);
configuration.setFetchTo(FetchTo.METADATA);
DeviceRelationsQuery deviceRelationsQuery = new DeviceRelationsQuery();
deviceRelationsQuery.setDirection(EntitySearchDirection.FROM);

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java

@ -16,15 +16,15 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
@Data
public class TbGetEntityAttrNodeConfiguration implements NodeConfiguration<TbGetEntityAttrNodeConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetEntityAttrNodeConfiguration> {
private Map<String, String> attrMapping;
private boolean isTelemetry = false;
@ -35,6 +35,7 @@ public class TbGetEntityAttrNodeConfiguration implements NodeConfiguration<TbGet
attrMapping.putIfAbsent("temperature", "tempo");
configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java

@ -16,14 +16,15 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.HashMap;
import java.util.Map;
@Data
public class TbGetOriginatorFieldsConfiguration implements NodeConfiguration<TbGetOriginatorFieldsConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetOriginatorFieldsConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetOriginatorFieldsConfiguration> {
private Map<String, String> fieldsMapping;
private boolean ignoreNullStrings;
@ -35,6 +36,7 @@ public class TbGetOriginatorFieldsConfiguration implements NodeConfiguration<TbG
fieldsMapping.put("type", "originatorType");
configuration.setFieldsMapping(fieldsMapping);
configuration.setIgnoreNullStrings(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

80
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java

@ -15,13 +15,13 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.node.ObjectNode;
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;
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;
@ -30,56 +30,78 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
/**
* Created by ashvayka on 19.01.18.
*/
@Slf4j
@RuleNode(type = ComponentType.ENRICHMENT,
name = "originator fields",
configClazz = TbGetOriginatorFieldsConfiguration.class,
nodeDescription = "Add Message Originator fields values into Message Metadata",
nodeDetails = "Will fetch fields values specified in mapping. If specified field is not part of originator fields it will be ignored.",
nodeDescription = "Add Message Originator fields values into Message Metadata or Message Data",
nodeDetails = "Will fetch fields values specified in mapping. If specified field is not part of originator fields it will be ignored. " +
"This node supports only following originator types: TENANT, CUSTOMER, USER, ASSET, DEVICE, ALARM, RULE_CHAIN, ENTITY_VIEW.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorFieldsConfig")
public class TbGetOriginatorFieldsNode implements TbNode {
private TbGetOriginatorFieldsConfiguration config;
private boolean ignoreNullStrings;
public class TbGetOriginatorFieldsNode extends TbAbstractNodeWithFetchTo<TbGetOriginatorFieldsConfiguration> {
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
config = TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class);
ignoreNullStrings = config.isIgnoreNullStrings();
protected TbGetOriginatorFieldsConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
try {
withCallback(putEntityFields(ctx, msg.getOriginator(), msg),
i -> ctx.tellSuccess(msg), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellFailure(msg, th);
ObjectNode msgDataAsJsonNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode = getMsgDataAsObjectNode(msg);
} else {
msgDataAsJsonNode = null;
}
ctx.checkTenantEntity(msg.getOriginator());
withCallback(collectMappedEntityFieldsAsync(ctx, msg.getOriginator()),
targetKeysToSourceValuesMap -> {
for (var entry : targetKeysToSourceValuesMap.entrySet()) {
var targetKeyName = entry.getKey();
var sourceFieldValue = entry.getValue();
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode.put(targetKeyName, sourceFieldValue);
} else if (FetchTo.METADATA.equals(fetchTo)) {
msg.getMetaData().putValue(targetKeyName, sourceFieldValue);
}
}
if (FetchTo.DATA.equals(fetchTo)) {
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataAsJsonNode)));
} else if (FetchTo.METADATA.equals(fetchTo)) {
ctx.tellSuccess(msg);
}
},
t -> ctx.tellFailure(msg, t),
MoreExecutors.directExecutor());
}
private ListenableFuture<Void> putEntityFields(TbContext ctx, EntityId entityId, TbMsg msg) {
private ListenableFuture<Map<String, String>> collectMappedEntityFieldsAsync(TbContext ctx, EntityId entityId) {
if (config.getFieldsMapping().isEmpty()) {
return Futures.immediateFuture(null);
return Futures.immediateFuture(Collections.emptyMap());
} else {
return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId),
data -> {
config.getFieldsMapping().forEach((field, metaKey) -> {
String val = data.getFieldValue(field, ignoreNullStrings);
if (val != null) {
msg.getMetaData().putValue(metaKey, val);
fieldsData -> {
var targetKeysToSourceValuesMap = new HashMap<String, String>();
for (var mappingEntry : config.getFieldsMapping().entrySet()) {
var sourceFieldName = mappingEntry.getKey();
var targetKeyName = mappingEntry.getValue();
var sourceFieldValue = fieldsData.getFieldValue(sourceFieldName, config.isIgnoreNullStrings());
if (sourceFieldValue != null) {
targetKeysToSourceValuesMap.put(targetKeyName, sourceFieldValue);
}
});
return null;
}, MoreExecutors.directExecutor()
}
return targetKeysToSourceValuesMap;
}, ctx.getDbCallbackExecutor()
);
}
}
}

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.data.RelationsQuery;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@ -26,8 +27,8 @@ import java.util.HashMap;
import java.util.Map;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfiguration {
private RelationsQuery relationsQuery;
@Override

25
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java

@ -27,30 +27,27 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
@RuleNode(
type = ComponentType.ENRICHMENT,
name="related attributes",
name = "related attributes",
configClazz = TbGetRelatedAttrNodeConfiguration.class,
nodeDescription = "Add Originators Related Entity Attributes or Latest Telemetry into Message Metadata",
nodeDescription = "Add Originators Related Entity Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "Related Entity found using configured relation direction and Relation Type. " +
"If multiple Related Entities are found, only first Entity is used for attributes enrichment, other entities are discarded. " +
"If Attributes enrichment configured, server scope attributes are added into Message metadata. " +
"If Latest Telemetry enrichment configured, latest telemetry added into metadata. " +
"If Attributes enrichment configured, server scope attributes are added into Message Metadata/Data. " +
"If Latest Telemetry enrichment configured, latest telemetry added into Metadata/Data. " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.temperature</code>.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeRelatedAttributesConfig")
public class TbGetRelatedAttributeNode extends TbEntityGetAttrNode<EntityId> {
private TbGetRelatedAttrNodeConfiguration config;
public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode<EntityId> {
@Override
public void init(TbContext context, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class);
setConfig(config);
public TbGetRelatedAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class);
}
@Override
protected ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, config.getRelationsQuery());
public ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
var relatedAttrConfig = (TbGetRelatedAttrNodeConfiguration) config;
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery());
}
}

18
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java

@ -20,6 +20,9 @@ import com.google.common.util.concurrent.ListenableFuture;
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.TbNodeConfiguration;
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.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -29,19 +32,22 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
type = ComponentType.ENRICHMENT,
name="tenant attributes",
configClazz = TbGetEntityAttrNodeConfiguration.class,
nodeDescription = "Add Originators Tenant Attributes or Latest Telemetry into Message Metadata",
nodeDetails = "If Attributes enrichment configured, server scope attributes are added into Message metadata. " +
"If Latest Telemetry enrichment configured, latest telemetry added into metadata. " +
nodeDescription = "Add Originators Tenant Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "If Attributes enrichment configured, server scope attributes are added into Message Metadata/Data. " +
"If Latest Telemetry enrichment configured, latest telemetry added into Metadata/Data. " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.temperature</code>.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeTenantAttributesConfig")
public class TbGetTenantAttributeNode extends TbEntityGetAttrNode<TenantId> {
public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode<TenantId> {
@Override
protected ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
public ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
return Futures.immediateFuture(ctx.getTenantId());
}
@Override
public TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
}

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

@ -25,6 +25,7 @@ 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.ContactBased;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -39,11 +40,10 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeEntityDetailsConfig")
public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode<TbGetTenantDetailsNodeConfiguration> {
private static final String TENANT_PREFIX = "tenant_";
@Override
protected TbGetTenantDetailsNodeConfiguration loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetTenantDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetTenantDetailsNodeConfiguration.class);
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java

@ -16,18 +16,19 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetTenantDetailsNodeConfiguration extends TbAbstractGetEntityDetailsNodeConfiguration implements NodeConfiguration<TbGetTenantDetailsNodeConfiguration> {
@Override
public TbGetTenantDetailsNodeConfiguration defaultConfiguration() {
TbGetTenantDetailsNodeConfiguration configuration = new TbGetTenantDetailsNodeConfiguration();
configuration.setDetailsList(Collections.emptyList());
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

42
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java

@ -1,12 +1,12 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* <p>
* 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
*
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.util;
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.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.BaseData;
@ -30,47 +31,50 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import java.util.function.Function;
@Slf4j
public class EntitiesFieldsAsyncLoader {
public static ListenableFuture<EntityFieldsData> findAsync(TbContext ctx, EntityId original) {
switch (original.getEntityType()) {
public static ListenableFuture<EntityFieldsData> findAsync(TbContext ctx, EntityId originatorId) {
switch (originatorId.getEntityType()) {
case TENANT:
return getAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), (TenantId) original),
return toEntityFieldsDataAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), (TenantId) originatorId),
EntityFieldsData::new);
case CUSTOMER:
return getAsync(ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), (CustomerId) original),
return toEntityFieldsDataAsync(ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), (CustomerId) originatorId),
EntityFieldsData::new);
case USER:
return getAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) original),
return toEntityFieldsDataAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originatorId),
EntityFieldsData::new);
case ASSET:
return getAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) original),
return toEntityFieldsDataAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originatorId),
EntityFieldsData::new);
case DEVICE:
return getAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) original),
return toEntityFieldsDataAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originatorId),
EntityFieldsData::new);
case ALARM:
return getAsync(ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), (AlarmId) original),
return toEntityFieldsDataAsync(ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), (AlarmId) originatorId),
EntityFieldsData::new);
case RULE_CHAIN:
return getAsync(ctx.getRuleChainService().findRuleChainByIdAsync(ctx.getTenantId(), (RuleChainId) original),
return toEntityFieldsDataAsync(ctx.getRuleChainService().findRuleChainByIdAsync(ctx.getTenantId(), (RuleChainId) originatorId),
EntityFieldsData::new);
case ENTITY_VIEW:
return getAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), (EntityViewId) original),
return toEntityFieldsDataAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), (EntityViewId) originatorId),
EntityFieldsData::new);
default:
return Futures.immediateFailedFuture(new TbNodeException("Unexpected original EntityType " + original.getEntityType()));
return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originatorId.getEntityType()));
}
}
private static <T extends BaseData> ListenableFuture<EntityFieldsData> getAsync(
ListenableFuture<T> future, Function<T, EntityFieldsData> converter) {
private static <T extends BaseData<? extends UUIDBased>> ListenableFuture<EntityFieldsData> toEntityFieldsDataAsync(
ListenableFuture<T> future,
Function<T, EntityFieldsData> converter
) {
return Futures.transformAsync(future, in -> in != null ?
Futures.immediateFuture(converter.apply(in))
: Futures.immediateFailedFuture(new RuntimeException("Entity not found!")), MoreExecutors.directExecutor());
: Futures.immediateFailedFuture(new TbNodeException("Entity not found!")), MoreExecutors.directExecutor());
}
}

20
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java

@ -1,12 +1,12 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* <p>
* 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
*
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
@ -65,7 +65,7 @@ import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public abstract class AbstractAttributeNodeTest {
public abstract class TbAbstractAttributeNodeTest {
final CustomerId customerId = new CustomerId(Uuids.timeBased());
final TenantId tenantId = TenantId.fromUUID(Uuids.timeBased());
final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased());
@ -86,9 +86,9 @@ public abstract class AbstractAttributeNodeTest {
DeviceService deviceService;
TbMsg msg;
Map<String, String> metaData;
TbEntityGetAttrNode node;
TbAbstractGetEntityAttrNode node;
void init(TbEntityGetAttrNode node) throws TbNodeException {
void init(TbAbstractGetEntityAttrNode node) throws TbNodeException {
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER;
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfig()));
@ -117,7 +117,6 @@ public abstract class AbstractAttributeNodeTest {
}
void errorThrownIfCannotLoadAttributesAsync(User user) {
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(ctx.getAttributesService()).thenReturn(attributesService);
@ -168,7 +167,7 @@ public abstract class AbstractAttributeNodeTest {
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER;
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigForTelemetry()));
TbEntityGetAttrNode node = getEmptyNode();
TbAbstractGetEntityAttrNode node = getEmptyNode();
node.init(null, nodeConfiguration);
msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
@ -210,10 +209,11 @@ public abstract class AbstractAttributeNodeTest {
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(isTelemetry);
config.setFetchTo(FetchTo.METADATA);
return config;
}
protected abstract TbEntityGetAttrNode getEmptyNode();
protected abstract TbAbstractGetEntityAttrNode getEmptyNode();
abstract EntityId getEntityId();

58
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java

@ -130,7 +130,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(false, false, false);
TbGetAttributesNode node = initNode(FetchTo.METADATA, false, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -138,9 +138,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, false, "cs_", clientAttributes);
checkAttributes(resultMsg, false, "ss_", serverAttributes);
checkAttributes(resultMsg, false, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries
checkTs(resultMsg, false, false, tsKeys);
@ -148,7 +148,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_latestWithTs_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(false, true, false);
TbGetAttributesNode node = initNode(FetchTo.METADATA, true, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -156,9 +156,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, false, "cs_", clientAttributes);
checkAttributes(resultMsg, false, "ss_", serverAttributes);
checkAttributes(resultMsg, false, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(resultMsg, false, true, tsKeys);
@ -166,7 +166,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(true, false, false);
TbGetAttributesNode node = initNode(FetchTo.DATA, false, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -174,9 +174,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, true, "cs_", clientAttributes);
checkAttributes(resultMsg, true, "ss_", serverAttributes);
checkAttributes(resultMsg, true, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries
checkTs(resultMsg, true, false, tsKeys);
@ -184,7 +184,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_latestWithTs_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(true, true, false);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -192,9 +192,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, true, "cs_", clientAttributes);
checkAttributes(resultMsg, true, "ss_", serverAttributes);
checkAttributes(resultMsg, true, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(resultMsg, true, true, tsKeys);
@ -202,7 +202,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_whenOnMsg_then_failure() throws Exception {
TbGetAttributesNode node = initNode(false, false, true);
TbGetAttributesNode node = initNode(FetchTo.METADATA, false, true);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -210,9 +210,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg actualMsg = checkMsg(false);
//check attributes
checkAttributes(actualMsg, false, "cs_", clientAttributes);
checkAttributes(actualMsg, false, "ss_", serverAttributes);
checkAttributes(actualMsg, false, "shared_", sharedAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(actualMsg, false, false, tsKeys);
@ -220,7 +220,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_then_failure() throws Exception {
TbGetAttributesNode node = initNode(true, true, true);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, true);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -228,9 +228,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg actualMsg = checkMsg(false);
//check attributes
checkAttributes(actualMsg, true, "cs_", clientAttributes);
checkAttributes(actualMsg, true, "ss_", serverAttributes);
checkAttributes(actualMsg, true, "shared_", sharedAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(actualMsg, true, true, tsKeys);
@ -238,7 +238,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_and_data_is_not_object_then_failure() throws Exception {
TbGetAttributesNode node = initNode(true, true, true);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, true);
TbMsg msg = TbMsg.newMsg("TEST", originator, new TbMsgMetaData(), "[]");
node.onMsg(ctx, msg);
@ -272,13 +272,13 @@ public class TbAbstractGetAttributesNodeTest {
return resultMsg;
}
private void checkAttributes(TbMsg actualMsg, boolean fetchToData, String prefix, List<String> attributes) {
private void checkAttributes(TbMsg actualMsg, FetchTo fetchTo, String prefix, List<String> attributes) {
JsonNode msgData = JacksonUtil.toJsonNode(actualMsg.getData());
attributes.stream()
.filter(attribute -> !attribute.equals("unknown"))
.forEach(attribute -> {
String result;
if (fetchToData) {
if (FetchTo.DATA.equals(fetchTo)) {
result = msgData.get(prefix + attribute).asText();
} else {
result = actualMsg.getMetaData().getValue(prefix + attribute);
@ -313,13 +313,13 @@ public class TbAbstractGetAttributesNodeTest {
}
}
private TbGetAttributesNode initNode(boolean fetchToData, boolean getLatestValueWithTs, boolean isTellFailureIfAbsent) throws TbNodeException {
private TbGetAttributesNode initNode(FetchTo fetchTo, boolean getLatestValueWithTs, boolean isTellFailureIfAbsent) throws TbNodeException {
TbGetAttributesNodeConfiguration config = new TbGetAttributesNodeConfiguration();
config.setClientAttributeNames(List.of("client_attr_1", "client_attr_2", "${client_attr_metadata}", "unknown"));
config.setServerAttributeNames(List.of("server_attr_1", "server_attr_2", "${server_attr_metadata}", "unknown"));
config.setSharedAttributeNames(List.of("shared_attr_1", "shared_attr_2", "$[shared_attr_data]", "unknown"));
config.setLatestTsKeyNames(List.of("temperature", "humidity", "unknown"));
config.setFetchToData(fetchToData);
config.setFetchTo(fetchTo);
config.setGetLatestValueWithTs(getLatestValueWithTs);
config.setTellFailureIfAbsent(isTellFailureIfAbsent);
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));

6
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java

@ -64,7 +64,7 @@ public class TbFetchDeviceCredentialsNodeTest {
callback = mock(TbMsgCallback.class);
ctx = mock(TbContext.class);
config = new TbFetchDeviceCredentialsNodeConfiguration().defaultConfiguration();
config.setFetchToMetadata(true);
config.setFetchTo(FetchTo.METADATA);
nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));
node = spy(new TbFetchDeviceCredentialsNode());
node.init(ctx, nodeConfiguration);
@ -89,13 +89,13 @@ public class TbFetchDeviceCredentialsNodeTest {
@Test
void givenDefaultConfig_whenInit_thenOK() {
assertThat(node.config).isEqualTo(config);
assertThat(node.fetchToMetadata).isEqualTo(true);
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA);
}
@Test
void givenDefaultConfig_whenVerify_thenOK() {
TbFetchDeviceCredentialsNodeConfiguration defaultConfig = new TbFetchDeviceCredentialsNodeConfiguration().defaultConfiguration();
assertThat(defaultConfig.isFetchToMetadata()).isEqualTo(true);
assertThat(defaultConfig.getFetchTo()).isEqualTo(FetchTo.METADATA);
}
@Test

81
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java

@ -1,12 +1,12 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* <p>
* 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
*
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
@ -16,28 +16,47 @@
package org.thingsboard.rule.engine.metadata;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.List;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetCustomerAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -56,7 +75,7 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetCustomerAttributeNode();
}
@ -65,6 +84,31 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
return customerId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetCustomerAttributeNode();
var config = new TbGetEntityAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new CustomerId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
mockFindUser(user);
@ -89,6 +133,29 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(customerId, "CUSTOMER");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
mockFindUser(user);

376
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java

@ -0,0 +1,376 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
* <p>
* 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
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.jetbrains.annotations.NotNull;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
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.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DashboardId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.device.DeviceService;
import java.util.Collections;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.Callable;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class TbGetOriginatorFieldsNodeTest {
private static final EntityId DUMMY_ENTITY_ID = new DeviceId(UUID.randomUUID());
public static final ListeningExecutor DB_EXECUTOR = new ListeningExecutor() {
@Override
public <T> ListenableFuture<T> executeAsync(Callable<T> task) {
try {
return Futures.immediateFuture(task.call());
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override
public void execute(@NotNull Runnable command) {
command.run();
}
};
@Mock
private TbContext ctxMock;
@Mock
private DeviceService deviceService;
private TbGetOriginatorFieldsNode node;
private TbGetOriginatorFieldsConfiguration config;
private TbNodeConfiguration nodeConfiguration;
private TbMsg msg;
@BeforeEach
public void setUp() {
config = new TbGetOriginatorFieldsConfiguration();
node = new TbGetOriginatorFieldsNode();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
}
@Test
public void givenConfigWithNullFetchTo_whenOnInit_thenException() {
// GIVEN
config = config.defaultConfiguration();
config.setFetchTo(null);
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration));
// THEN
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException {
// GIVEN
config = config.defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
node.init(ctxMock, nodeConfiguration);
// THEN
assertThat(node.config).isEqualTo(config);
assertThat(config.getFieldsMapping()).isEqualTo(Map.of(
"name", "originatorName",
"type", "originatorType"));
assertThat(config.isIgnoreNullStrings()).isEqualTo(false);
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA);
}
@Test
public void givenCustomConfig_whenInit_thenOK() throws TbNodeException {
// GIVEN
config.setFieldsMapping(Map.of(
"sourceField1", "targetKey1",
"sourceField2", "targetKey2",
"sourceField3", "targetKey3"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.DATA);
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
node.init(ctxMock, nodeConfiguration);
// THEN
assertThat(node.config).isEqualTo(config);
assertThat(config.getFieldsMapping()).isEqualTo(Map.of(
"sourceField1", "targetKey1",
"sourceField2", "targetKey2",
"sourceField3", "targetKey3"));
assertThat(config.isIgnoreNullStrings()).isEqualTo(true);
assertThat(node.fetchTo).isEqualTo(FetchTo.DATA);
}
@Test
public void givenMsgDataIsNotAnJsonObjectAndFetchToData_whenOnMsg_thenException() {
// GIVEN
node.fetchTo = FetchTo.DATA;
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_ENTITY_ID, new TbMsgMetaData(), "[]");
// WHEN
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctxMock, msg));
// THEN
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenEntityThatDoesNotBelongToTheCurrentTenant_whenOnMsg_thenException() {
// SETUP
var expectedExceptionMessage = "Entity with id: '" + DUMMY_ENTITY_ID +
"' specified in the configuration doesn't belong to the current tenant.";
// GIVEN
doThrow(new RuntimeException(expectedExceptionMessage)).when(ctxMock).checkTenantEntity(DUMMY_ENTITY_ID);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_ENTITY_ID, new TbMsgMetaData(), "{}");
// WHEN
var exception = assertThrows(RuntimeException.class, () -> node.onMsg(ctxMock, msg));
// THEN
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMessage);
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenValidMsgAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchToData() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.DATA);
node.config = config;
node.fetchTo = FetchTo.DATA;
var msgMetaData = new TbMsgMetaData();
var msgData = "{\"temp\":42,\"humidity\":77}";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgData = "{\"temp\":42,\"humidity\":77,\"originatorName\":\"Test device\",\"originatorType\":\"Test device type\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msgMetaData);
}
@Test
public void givenValidMsgAndFetchToMetaData_whenOnMsg_thenShouldTellSuccessAndFetchToMetaData() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123",
"originatorName", "Test device",
"originatorType", "Test device type"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenNullEntityFieldsAndIgnoreNullStringsFalse_whenOnMsg_thenShouldTellSuccessAndFetchNullField() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123",
"originatorName", "Test device",
"originatorType", "Test device type",
"originatorLabel", "null"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenEmptyFieldsMapping_whenOnMsg_thenShouldTellSuccessWithSameMsg() {
// GIVEN
config.setFieldsMapping(Collections.emptyMap());
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
var expectedMsgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenUnsupportedEntityType_whenOnMsg_thenShouldTellFailureWithSameMsg() {
// GIVEN
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", new DashboardId(UUID.randomUUID()), msgMetaData, msgData);
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellFailure(actualMessageCaptor.capture(), any());
verify(ctxMock, never()).tellSuccess(any());
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msgMetaData);
}
}

83
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java

@ -1,12 +1,12 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* <p>
* 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
*
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
@ -15,12 +15,16 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
@ -29,7 +33,13 @@ import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
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.dao.relation.RelationService;
import java.util.HashMap;
@ -37,11 +47,19 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetRelatedAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -69,7 +87,7 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetRelatedAttributeNode();
}
@ -90,6 +108,7 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(isTelemetry);
config.setFetchTo(FetchTo.METADATA);
return config;
}
@ -98,6 +117,31 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
return customerId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetRelatedAttributeNode();
var config = new TbGetRelatedAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new DeviceId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
entityRelation.setFrom(user.getId());
@ -130,6 +174,33 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(customerId, "CUSTOMER");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
entityRelation.setFrom(customerId);
entityRelation.setTo(customerId);
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation)));
msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
entityRelation.setFrom(user.getId());

85
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java

@ -1,12 +1,12 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* <p>
* 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
*
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
@ -15,10 +15,15 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
@ -26,15 +31,31 @@ import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.List;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetTenantAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -56,7 +77,7 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetTenantAttributeNode();
}
@ -65,6 +86,31 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
return tenantId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetTenantAttributeNode();
var config = new TbGetEntityAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new TenantId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
errorThrownIfCannotLoadAttributes(user);
@ -86,6 +132,29 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(tenantId, "TENANT");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("TENANT", tenantId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
usersCustomerAttributesFetched(user);

231
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java

@ -0,0 +1,231 @@
package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityFieldsData;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.user.UserService;
import java.util.EnumSet;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class EntitiesFieldsAsyncLoaderTest {
private static EnumSet<EntityType> SUPPORTED_ENTITY_TYPES;
private static UUID RANDOM_UUID;
private static TenantId TENANT_ID;
@Mock
private TbContext ctxMock;
@Mock
private TenantService tenantServiceMock;
@Mock
private CustomerService customerServiceMock;
@Mock
private UserService userServiceMock;
@Mock
private AssetService assetServiceMock;
@Mock
private DeviceService deviceServiceMock;
@Mock
private RuleEngineAlarmService ruleEngineAlarmServiceMock;
@Mock
private RuleChainService ruleChainServiceMock;
@Mock
private EntityViewService entityViewServiceMock;
@BeforeAll
public static void setup() {
RANDOM_UUID = UUID.randomUUID();
TENANT_ID = new TenantId(UUID.randomUUID());
SUPPORTED_ENTITY_TYPES = EnumSet.of(
EntityType.TENANT,
EntityType.CUSTOMER,
EntityType.USER,
EntityType.ASSET,
EntityType.DEVICE,
EntityType.ALARM,
EntityType.RULE_CHAIN,
EntityType.ENTITY_VIEW
);
}
@Test
public void givenSupportedEntityTypes_whenFindAsync_thenOK() throws ExecutionException, InterruptedException {
for (var entityType : SUPPORTED_ENTITY_TYPES) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
initMocks(entityType, false);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
var actualEntityFieldsData = EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get();
var expectedEntityFieldsData = new EntityFieldsData(getEntityFromEntityId(entityId));
Assertions.assertEquals(expectedEntityFieldsData, actualEntityFieldsData);
}
}
@Test
public void givenUnsupportedEntityTypes_whenFindAsync_thenException() {
for (var entityType : EntityType.values()) {
if (!SUPPORTED_ENTITY_TYPES.contains(entityType)) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Unexpected originator EntityType: " + entityType;
var exception = assertThrows(ExecutionException.class,
() -> EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get());
assertInstanceOf(TbNodeException.class, exception.getCause());
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMsg);
}
}
}
@Test
public void givenSupportedTypeButEntityDoesNotExist_whenFindAsync_thenException() {
for (var entityType : SUPPORTED_ENTITY_TYPES) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
initMocks(entityType, true);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Entity not found!";
var exception = assertThrows(ExecutionException.class,
() -> EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get());
assertInstanceOf(TbNodeException.class, exception.getCause());
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMsg);
}
}
private void initMocks(EntityType entityType, boolean entityDoesNotExist) {
switch (entityType) {
case TENANT:
var tenant = Futures.immediateFuture(entityDoesNotExist ? null : new Tenant(new TenantId(RANDOM_UUID)));
when(ctxMock.getTenantService()).thenReturn(tenantServiceMock);
doReturn(tenant).when(tenantServiceMock).findTenantByIdAsync(eq(TENANT_ID), any());
break;
case CUSTOMER:
var customer = Futures.immediateFuture(entityDoesNotExist ? null : new Customer(new CustomerId(RANDOM_UUID)));
when(ctxMock.getCustomerService()).thenReturn(customerServiceMock);
doReturn(customer).when(customerServiceMock).findCustomerByIdAsync(eq(TENANT_ID), any());
break;
case USER:
var user = Futures.immediateFuture(entityDoesNotExist ? null : new User(new UserId(RANDOM_UUID)));
when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(user).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), any());
break;
case ASSET:
var asset = Futures.immediateFuture(entityDoesNotExist ? null : new Asset(new AssetId(RANDOM_UUID)));
when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
doReturn(asset).when(assetServiceMock).findAssetByIdAsync(eq(TENANT_ID), any());
break;
case DEVICE:
var device = Futures.immediateFuture(entityDoesNotExist ? null : new Device(new DeviceId(RANDOM_UUID)));
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
doReturn(device).when(deviceServiceMock).findDeviceByIdAsync(eq(TENANT_ID), any());
break;
case ALARM:
var alarm = Futures.immediateFuture(entityDoesNotExist ? null : new Alarm(new AlarmId(RANDOM_UUID)));
when(ctxMock.getAlarmService()).thenReturn(ruleEngineAlarmServiceMock);
doReturn(alarm).when(ruleEngineAlarmServiceMock).findAlarmByIdAsync(eq(TENANT_ID), any());
break;
case RULE_CHAIN:
var ruleChain = Futures.immediateFuture(entityDoesNotExist ? null : new RuleChain(new RuleChainId(RANDOM_UUID)));
when(ctxMock.getRuleChainService()).thenReturn(ruleChainServiceMock);
doReturn(ruleChain).when(ruleChainServiceMock).findRuleChainByIdAsync(eq(TENANT_ID), any());
break;
case ENTITY_VIEW:
var entityView = Futures.immediateFuture(entityDoesNotExist ? null : new EntityView(new EntityViewId(RANDOM_UUID)));
when(ctxMock.getEntityViewService()).thenReturn(entityViewServiceMock);
doReturn(entityView).when(entityViewServiceMock).findEntityViewByIdAsync(eq(TENANT_ID), any());
break;
default:
throw new RuntimeException("Unexpected EntityType: " + entityType);
}
}
private BaseData<? extends UUIDBased> getEntityFromEntityId(EntityId entityId) {
switch (entityId.getEntityType()) {
case TENANT:
return new Tenant((TenantId) entityId);
case CUSTOMER:
return new Customer((CustomerId) entityId);
case USER:
return new User((UserId) entityId);
case ASSET:
return new Asset((AssetId) entityId);
case DEVICE:
return new Device((DeviceId) entityId);
case ALARM:
return new Alarm((AlarmId) entityId);
case RULE_CHAIN:
return new RuleChain((RuleChainId) entityId);
case ENTITY_VIEW:
return new EntityView((EntityViewId) entityId);
default:
throw new RuntimeException("Unexpected EntityType: " + entityId.getEntityType());
}
}
}

7
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java

@ -56,7 +56,6 @@ import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.dao.alarm.AlarmCommentService;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.dashboard.DashboardService;
@ -94,8 +93,6 @@ public class TenantIdLoaderTest {
@Mock
private RuleEngineAlarmService alarmService;
@Mock
private AlarmCommentService alarmCommentService;
@Mock
private RuleChainService ruleChainService;
@Mock
private EntityViewService entityViewService;
@ -312,9 +309,8 @@ public class TenantIdLoaderTest {
break;
default:
throw new RuntimeException("Unexpected original EntityType " + entityType);
throw new RuntimeException("Unexpected originator EntityType " + entityType);
}
}
private EntityId getEntityId(EntityType entityType) {
@ -350,5 +346,4 @@ public class TenantIdLoaderTest {
public void test_findEntityIdAsync_other_tenant() {
checkTenant(new TenantId(UUID.randomUUID()), false);
}
}

Loading…
Cancel
Save