Browse Source

enrichment nodes improvements

pull/8661/head
ShvaykaD 4 years ago
parent
commit
30b7d819ec
  1. 33
      common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java
  2. 31
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java
  3. 94
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java
  4. 1
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java
  5. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java

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

@ -19,5 +19,36 @@ package org.thingsboard.server.common.data;
* @author Andrew Shvayka
*/
public enum EntityType {
TENANT, CUSTOMER, USER, DASHBOARD, ASSET, DEVICE, ALARM, RULE_CHAIN, RULE_NODE, ENTITY_VIEW, WIDGETS_BUNDLE, WIDGET_TYPE, TENANT_PROFILE, DEVICE_PROFILE, ASSET_PROFILE, API_USAGE_STATE, TB_RESOURCE, OTA_PACKAGE, EDGE, RPC, QUEUE;
TENANT("Tenant"),
CUSTOMER("Customer"),
USER("User"),
DASHBOARD("Dashboard"),
ASSET("Asset"),
DEVICE("Device"),
ALARM("Alarm"),
RULE_CHAIN("Rule chain"),
RULE_NODE("Rule node"),
ENTITY_VIEW("Entity view"),
WIDGETS_BUNDLE("Widget bundle"),
WIDGET_TYPE("Widget type"),
TENANT_PROFILE("Tenant profile"),
DEVICE_PROFILE("Device profile"),
ASSET_PROFILE("Asset profile"),
API_USAGE_STATE("Api usage state"),
TB_RESOURCE("TB resource"),
OTA_PACKAGE("OTA package"),
EDGE("Edge"),
RPC("Rpc"),
QUEUE("Queue");
private final String displayName;
EntityType(String displayName) {
this.displayName = displayName;
}
public String getDisplayName() {
return displayName;
}
}

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

@ -39,7 +39,6 @@ import java.lang.reflect.Type;
import java.util.Map;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@Slf4j
public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEntityDetailsNodeConfiguration> implements TbNode {
@ -71,47 +70,47 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
protected MessageData getDataAsJson(TbMsg msg) {
if (this.config.isAddToMetadata()) {
return new MessageData(gson.toJsonTree(msg.getMetaData().getData(), TYPE), "metadata");
return new MessageData(gson.toJsonTree(msg.getMetaData().getData(), TYPE), DataSource.METADATA);
} else {
return new MessageData(jsonParser.parse(msg.getData()), "data");
return new MessageData(jsonParser.parse(msg.getData()), DataSource.DATA);
}
}
protected ListenableFuture<TbMsg> getTbMsgListenableFuture(TbContext ctx, TbMsg msg, MessageData messageData, String prefix) {
if (!this.config.getDetailsList().isEmpty()) {
if (this.config.getDetailsList().isEmpty()) {
return Futures.immediateFuture(msg);
} else {
ListenableFuture<ContactBased> contactBasedListenableFuture = getContactBasedListenableFuture(ctx, msg);
ListenableFuture<JsonElement> resultObject = addContactProperties(messageData.getData(), contactBasedListenableFuture, prefix);
return transformMsg(ctx, msg, resultObject, messageData);
} else {
return Futures.immediateFuture(msg);
}
}
private ListenableFuture<TbMsg> transformMsg(TbContext ctx, TbMsg msg, ListenableFuture<JsonElement> propertiesFuture, MessageData messageData) {
return Futures.transformAsync(propertiesFuture, jsonElement -> {
if (jsonElement != null) {
if (messageData.getDataType().equals("metadata")) {
if (jsonElement == null) {
return Futures.immediateFuture(null);
} else {
if (messageData.getDataSource().equals(DataSource.METADATA)) {
Map<String, String> metadataMap = gson.fromJson(jsonElement.toString(), TYPE);
return Futures.immediateFuture(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), new TbMsgMetaData(metadataMap), msg.getData()));
} else {
return Futures.immediateFuture(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), gson.toJson(jsonElement)));
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
}
private ListenableFuture<JsonElement> addContactProperties(JsonElement data, ListenableFuture<ContactBased> entityFuture, String prefix) {
return Futures.transformAsync(entityFuture, contactBased -> {
if (contactBased != null) {
if (contactBased == null) {
return Futures.immediateFuture(null);
} else {
JsonElement jsonElement = null;
for (EntityDetails entityDetails : this.config.getDetailsList()) {
jsonElement = setProperties(contactBased, data, entityDetails, prefix);
}
return Futures.immediateFuture(jsonElement);
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
}
@ -175,7 +174,11 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
@AllArgsConstructor
private static class MessageData {
private JsonElement data;
private String dataType;
private DataSource dataSource;
}
private enum DataSource {
DATA, METADATA
}

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

@ -26,9 +26,12 @@ 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.Customer;
import org.thingsboard.server.common.data.HasCustomerId;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.id.AssetId;
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.UserId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -60,80 +63,47 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode<TbG
@Override
protected ListenableFuture<ContactBased> getContactBasedListenableFuture(TbContext ctx, TbMsg msg) {
return Futures.transformAsync(getCustomer(ctx, msg), customer -> {
if (customer != null) {
return Futures.immediateFuture(customer);
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(getCustomer(ctx, msg), customer ->
customer == null ? Futures.immediateFuture(null) : Futures.immediateFuture(customer),
MoreExecutors.directExecutor());
}
private ListenableFuture<Customer> getCustomer(TbContext ctx, TbMsg msg) {
switch (msg.getOriginator().getEntityType()) {
case DEVICE:
return Futures.transformAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), new DeviceId(msg.getOriginator().getId())), device -> {
if (device != null) {
if (!device.getCustomerId().isNullUid()) {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), device.getCustomerId());
} else {
throw new RuntimeException("Device with name '" + device.getName() + "' is not assigned to Customer.");
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), new DeviceId(msg.getOriginator().getId())),
device -> getCustomerFuture(ctx, device, msg.getOriginator()), MoreExecutors.directExecutor());
case ASSET:
return Futures.transformAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), new AssetId(msg.getOriginator().getId())), asset -> {
if (asset != null) {
if (!asset.getCustomerId().isNullUid()) {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), asset.getCustomerId());
} else {
throw new RuntimeException("Asset with name '" + asset.getName() + "' is not assigned to Customer.");
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), new AssetId(msg.getOriginator().getId())),
asset -> getCustomerFuture(ctx, asset, msg.getOriginator()), MoreExecutors.directExecutor());
case ENTITY_VIEW:
return Futures.transformAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), new EntityViewId(msg.getOriginator().getId())), entityView -> {
if (entityView != null) {
if (!entityView.getCustomerId().isNullUid()) {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), entityView.getCustomerId());
} else {
throw new RuntimeException("EntityView with name '" + entityView.getName() + "' is not assigned to Customer.");
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), new EntityViewId(msg.getOriginator().getId())),
entityView -> getCustomerFuture(ctx, entityView, msg.getOriginator()), MoreExecutors.directExecutor());
case USER:
return Futures.transformAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), new UserId(msg.getOriginator().getId())), user -> {
if (user != null) {
if (!user.getCustomerId().isNullUid()) {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), user.getCustomerId());
} else {
throw new RuntimeException("User with name '" + user.getName() + "' is not assigned to Customer.");
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), new UserId(msg.getOriginator().getId())),
user -> getCustomerFuture(ctx, user, msg.getOriginator()), MoreExecutors.directExecutor());
case EDGE:
return Futures.transformAsync(ctx.getEdgeService().findEdgeByIdAsync(ctx.getTenantId(), new EdgeId(msg.getOriginator().getId())), edge -> {
if (edge != null) {
if (!edge.getCustomerId().isNullUid()) {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), edge.getCustomerId());
} else {
throw new RuntimeException("Edge with name '" + edge.getName() + "' is not assigned to Customer.");
}
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
return Futures.transformAsync(ctx.getEdgeService().findEdgeByIdAsync(ctx.getTenantId(), new EdgeId(msg.getOriginator().getId())),
edge -> getCustomerFuture(ctx, edge, msg.getOriginator()), MoreExecutors.directExecutor());
default:
throw new RuntimeException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported.");
}
}
private ListenableFuture<Customer> getCustomerFuture(TbContext ctx, HasCustomerId hasCustomerId, EntityId originator) {
if (hasCustomerId == null) {
return Futures.immediateFuture(null);
} else {
if (hasCustomerId.getCustomerId().isNullUid()) {
if (hasCustomerId instanceof HasName) {
HasName hasName = (HasName) hasCustomerId;
throw new RuntimeException(originator.getEntityType().getDisplayName() + " with name '" + hasName.getName() + "' is not assigned to Customer.");
}
throw new RuntimeException(originator.getEntityType().getDisplayName() + " with id '" + originator + "' is not assigned to Customer.");
} else {
return ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), hasCustomerId.getCustomerId());
}
}
}
}

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

@ -40,6 +40,7 @@ public class TbGetTenantAttributeNode extends TbEntityGetAttrNode<TenantId> {
@Override
protected ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
return Futures.immediateFuture(ctx.getTenantId());
}

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

@ -54,12 +54,9 @@ public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode<TbGet
@Override
protected ListenableFuture<ContactBased> getContactBasedListenableFuture(TbContext ctx, TbMsg msg) {
return Futures.transformAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), ctx.getTenantId()), tenant -> {
if (tenant != null) {
return Futures.immediateFuture(tenant);
} else {
return Futures.immediateFuture(null);
}
}, MoreExecutors.directExecutor());
ctx.checkTenantEntity(msg.getOriginator());
return Futures.transformAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), ctx.getTenantId()), tenant ->
tenant == null ? Futures.immediateFuture(null) : Futures.immediateFuture(tenant),
MoreExecutors.directExecutor());
}
}

Loading…
Cancel
Save