Browse Source

Merge pull request #7245 from YuriyLytvynchuk/feature/node_change_originator_add_entitytype

[3.4.2] Feature: change originator node - add ENTITY_SOURCE
pull/7341/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
f229a506cc
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java
  2. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java
  3. 39
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java
  4. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java
  5. 90
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesByNameAndTypeLoader.java
  6. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesTenantIdAsyncLoader.java
  7. 17
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java

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

@ -40,8 +40,6 @@ import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
@ -253,11 +251,9 @@ public abstract class TbAbstractRelationActionNode<C extends TbAbstractRelationA
break;
case DASHBOARD:
DashboardService dashboardService = ctx.getDashboardService();
PageData<DashboardInfo> dashboardInfoTextPageData = dashboardService.findDashboardsByTenantId(ctx.getTenantId(), new PageLink(200, 0, entitykey.getEntityName()));
for (DashboardInfo dashboardInfo : dashboardInfoTextPageData.getData()) {
if (dashboardInfo.getTitle().equals(entitykey.getEntityName())) {
targetEntity.setEntityId(dashboardInfo.getId());
}
DashboardInfo dashboardInfo = dashboardService.findFirstDashboardInfoByTenantIdAndName(ctx.getTenantId(), entitykey.getEntityName());
if (dashboardInfo != null) {
targetEntity.setEntityId(dashboardInfo.getId());
}
break;
case USER:

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

@ -15,11 +15,11 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures;
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.util.EntitiesTenantIdAsyncLoader;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -40,7 +40,7 @@ public class TbGetTenantAttributeNode extends TbEntityGetAttrNode<TenantId> {
@Override
protected ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
return EntitiesTenantIdAsyncLoader.findEntityIdAsync(ctx, originator);
return Futures.immediateFuture(ctx.getTenantId());
}
}

39
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java

@ -25,9 +25,11 @@ 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.EntitiesAlarmOriginatorIdAsyncLoader;
import org.thingsboard.rule.engine.util.EntitiesByNameAndTypeLoader;
import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader;
import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader;
import org.thingsboard.rule.engine.util.EntitiesTenantIdAsyncLoader;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
@ -55,6 +57,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode {
protected static final String TENANT_SOURCE = "TENANT";
protected static final String RELATED_SOURCE = "RELATED";
protected static final String ALARM_ORIGINATOR_SOURCE = "ALARM_ORIGINATOR";
protected static final String ENTITY_SOURCE = "ENTITY_SOURCE";
private TbChangeOriginatorNodeConfiguration config;
@ -67,7 +70,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode {
@Override
protected ListenableFuture<List<TbMsg>> transform(TbContext ctx, TbMsg msg) {
ListenableFuture<? extends EntityId> newOriginator = getNewOriginator(ctx, msg.getOriginator());
ListenableFuture<? extends EntityId> newOriginator = getNewOriginator(ctx, msg);
return Futures.transform(newOriginator, n -> {
if (n == null || n.isNullUid()) {
return null;
@ -76,23 +79,32 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode {
}, ctx.getDbCallbackExecutor());
}
private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, EntityId original) {
private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, TbMsg msg) {
switch (config.getOriginatorSource()) {
case CUSTOMER_SOURCE:
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, original);
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator());
case TENANT_SOURCE:
return EntitiesTenantIdAsyncLoader.findEntityIdAsync(ctx, original);
return Futures.immediateFuture(ctx.getTenantId());
case RELATED_SOURCE:
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, original, config.getRelationsQuery());
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery());
case ALARM_ORIGINATOR_SOURCE:
return EntitiesAlarmOriginatorIdAsyncLoader.findEntityIdAsync(ctx, original);
return EntitiesAlarmOriginatorIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator());
case ENTITY_SOURCE:
EntityType entityType = EntityType.valueOf(config.getEntityType());
String entityName = TbNodeUtils.processPattern(config.getEntityNamePattern(), msg);
try {
EntityId targetEntity = EntitiesByNameAndTypeLoader.findEntityId(ctx, entityType, entityName);
return Futures.immediateFuture(targetEntity);
} catch (IllegalStateException e) {
return Futures.immediateFailedFuture(e);
}
default:
return Futures.immediateFailedFuture(new IllegalStateException("Unexpected originator source " + config.getOriginatorSource()));
}
}
private void validateConfig(TbChangeOriginatorNodeConfiguration conf) {
HashSet<String> knownSources = Sets.newHashSet(CUSTOMER_SOURCE, TENANT_SOURCE, RELATED_SOURCE, ALARM_ORIGINATOR_SOURCE);
HashSet<String> knownSources = Sets.newHashSet(CUSTOMER_SOURCE, TENANT_SOURCE, RELATED_SOURCE, ALARM_ORIGINATOR_SOURCE, ENTITY_SOURCE);
if (!knownSources.contains(conf.getOriginatorSource())) {
log.error("Unsupported source [{}] for TbChangeOriginatorNode", conf.getOriginatorSource());
throw new IllegalArgumentException("Unsupported source TbChangeOriginatorNode" + conf.getOriginatorSource());
@ -106,6 +118,17 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode {
}
}
if (conf.getOriginatorSource().equals(ENTITY_SOURCE)) {
if (conf.getEntityType() == null) {
log.error("Entity type not specified for [{}]", ENTITY_SOURCE);
throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE);
}
if (StringUtils.isEmpty(conf.getEntityNamePattern())) {
log.error("EntityNamePattern not specified for type [{}]", conf.getEntityType());
throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE);
}
}
}
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java

@ -30,6 +30,8 @@ public class TbChangeOriginatorNodeConfiguration extends TbTransformNodeConfigur
private String originatorSource;
private RelationsQuery relationsQuery;
private String entityType;
private String entityNamePattern;
@Override
public TbChangeOriginatorNodeConfiguration defaultConfiguration() {

90
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesByNameAndTypeLoader.java

@ -0,0 +1,90 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.util;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DashboardInfo;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EntityId;
import java.util.Optional;
public class EntitiesByNameAndTypeLoader {
public static EntityId findEntityId(TbContext ctx, EntityType entityType, String entityName) {
EntityId targetEntityId = null;
switch (entityType) {
case DEVICE:
Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(ctx.getTenantId(), entityName);
if (device != null) {
targetEntityId = device.getId();
}
break;
case ASSET:
Asset asset = ctx.getAssetService().findAssetByTenantIdAndName(ctx.getTenantId(), entityName);
if (asset != null) {
targetEntityId = asset.getId();
}
break;
case CUSTOMER:
Optional<Customer> customerOptional = ctx.getCustomerService().findCustomerByTenantIdAndTitle(ctx.getTenantId(), entityName);
if (customerOptional.isPresent()) {
targetEntityId = customerOptional.get().getId();
}
break;
case TENANT:
targetEntityId = ctx.getTenantId();
break;
case ENTITY_VIEW:
EntityView entityView = ctx.getEntityViewService().findEntityViewByTenantIdAndName(ctx.getTenantId(), entityName);
if (entityView != null) {
targetEntityId = entityView.getId();
}
break;
case EDGE:
Edge edge = ctx.getEdgeService().findEdgeByTenantIdAndName(ctx.getTenantId(), entityName);
if (edge != null) {
targetEntityId = edge.getId();
}
break;
case DASHBOARD:
DashboardInfo dashboardInfo = ctx.getDashboardService().findFirstDashboardInfoByTenantIdAndName(ctx.getTenantId(), entityName);
if (dashboardInfo != null) {
targetEntityId = dashboardInfo.getId();
}
break;
case USER:
User user = ctx.getUserService().findUserByEmail(ctx.getTenantId(), entityName);
if (user != null) {
targetEntityId = user.getId();
}
break;
default:
throw new IllegalStateException("Unexpected entity type " + entityType.name());
}
if (targetEntityId == null) {
throw new IllegalStateException("Failed to found " + entityType.name() + " entity by name: '" + entityName + "'!");
}
return targetEntityId;
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesTenantIdAsyncLoader.java

@ -31,7 +31,10 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
public class EntitiesTenantIdAsyncLoader {
/**
* @deprecated consider to remove since tenantId is already defined in the TbContext.
*/
@Deprecated
public static ListenableFuture<TenantId> findEntityIdAsync(TbContext ctx, EntityId original) {
switch (original.getEntityType()) {

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

@ -15,7 +15,6 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
@ -31,8 +30,6 @@ import org.thingsboard.server.common.data.id.UserId;
import java.util.UUID;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.when;
@RunWith(MockitoJUnitRunner.class)
@ -45,6 +42,7 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
@Before
public void initDataForTests() throws TbNodeException {
init(new TbGetTenantAttributeNode());
user.setTenantId(tenantId);
user.setId(new UserId(UUID.randomUUID()));
@ -53,6 +51,8 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
device.setTenantId(tenantId);
device.setId(new DeviceId(UUID.randomUUID()));
when(ctx.getTenantId()).thenReturn(tenantId);
}
@Override
@ -67,20 +67,17 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
@Test
public void errorThrownIfCannotLoadAttributes() {
mockFindUser(user);
errorThrownIfCannotLoadAttributes(user);
}
@Test
public void errorThrownIfCannotLoadAttributesAsync() {
mockFindUser(user);
errorThrownIfCannotLoadAttributesAsync(user);
}
@Test
public void failedChainUsedIfCustomerCannotBeFound() {
when(ctx.getUserService()).thenReturn(userService);
when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null));
public void failedChainUsedIfTenantIdFromCtxCannotBeFound() {
when(ctx.getTenantId()).thenReturn(null);
failedChainUsedIfCustomerCannotBeFound(user);
}
@ -91,25 +88,21 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
@Test
public void usersCustomerAttributesFetched() {
mockFindUser(user);
usersCustomerAttributesFetched(user);
}
@Test
public void assetsCustomerAttributesFetched() {
mockFindAsset(asset);
assetsCustomerAttributesFetched(asset);
}
@Test
public void deviceCustomerAttributesFetched() {
mockFindDevice(device);
deviceCustomerAttributesFetched(device);
}
@Test
public void deviceCustomerTelemetryFetched() throws TbNodeException {
mockFindDevice(device);
deviceCustomerTelemetryFetched(device);
}
}

Loading…
Cancel
Save