diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/OriginatorSource.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/OriginatorSource.java new file mode 100644 index 0000000000..2ae5402404 --- /dev/null +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/OriginatorSource.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2024 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.transform; + +public enum OriginatorSource { + CUSTOMER, + TENANT, + RELATED, + ALARM_ORIGINATOR, + ENTITY +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java index 7bf6338e81..d590bbbb75 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java @@ -15,7 +15,6 @@ */ package org.thingsboard.rule.engine.transform; -import com.google.common.collect.Sets; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -34,10 +33,12 @@ 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.HashSet; import java.util.List; import java.util.NoSuchElementException; +import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY; +import static org.thingsboard.rule.engine.transform.OriginatorSource.RELATED; + @Slf4j @RuleNode( type = ComponentType.TRANSFORMATION, @@ -59,12 +60,6 @@ import java.util.NoSuchElementException; ) public class TbChangeOriginatorNode extends TbAbstractTransformNode { - private static final String CUSTOMER_SOURCE = "CUSTOMER"; - private static final String TENANT_SOURCE = "TENANT"; - private static final String RELATED_SOURCE = "RELATED"; - private static final String ALARM_ORIGINATOR_SOURCE = "ALARM_ORIGINATOR"; - private static final String ENTITY_SOURCE = "ENTITY"; - @Override protected TbChangeOriginatorNodeConfiguration loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { var config = TbNodeUtils.convert(configuration, TbChangeOriginatorNodeConfiguration.class); @@ -85,15 +80,15 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode getNewOriginator(TbContext ctx, TbMsg msg) { switch (config.getOriginatorSource()) { - case CUSTOMER_SOURCE: + case CUSTOMER: return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator()); - case TENANT_SOURCE: + case TENANT: return Futures.immediateFuture(ctx.getTenantId()); - case RELATED_SOURCE: + case RELATED: return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery()); - case ALARM_ORIGINATOR_SOURCE: + case ALARM_ORIGINATOR: return EntitiesAlarmOriginatorIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator()); - case ENTITY_SOURCE: + case ENTITY: EntityType entityType = EntityType.valueOf(config.getEntityType()); String entityName = TbNodeUtils.processPattern(config.getEntityNamePattern(), msg); try { @@ -108,28 +103,22 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode 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()); + if (conf.getOriginatorSource() == null) { + log.debug("Originator source should be specified."); + throw new IllegalArgumentException("Originator source should be specified."); } - - if (conf.getOriginatorSource().equals(RELATED_SOURCE)) { - if (conf.getRelationsQuery() == null) { - log.error("Related source for TbChangeOriginatorNode should have relations query. Actual [{}]", - conf.getRelationsQuery()); - throw new IllegalArgumentException("Wrong config for RElated Source in TbChangeOriginatorNode" + conf.getOriginatorSource()); - } + if (conf.getOriginatorSource().equals(RELATED) && conf.getRelationsQuery() == null) { + log.debug("Relations query should be specified if 'Related entity' source is selected."); + throw new IllegalArgumentException("Relations query should be specified if 'Related entity' source is selected."); } - - if (conf.getOriginatorSource().equals(ENTITY_SOURCE)) { + if (conf.getOriginatorSource().equals(ENTITY)) { if (conf.getEntityType() == null) { - log.error("Entity type not specified for [{}]", ENTITY_SOURCE); - throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE); + log.debug("Entity type should be specified if '{}' source is selected.", ENTITY); + throw new IllegalArgumentException("Entity type should be specified if 'Entity by name pattern' source is selected."); } if (StringUtils.isEmpty(conf.getEntityNamePattern())) { - log.error("EntityNamePattern not specified for type [{}]", conf.getEntityType()); - throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE); + log.debug("Name pattern should be specified if '{}' source is selected.", ENTITY); + throw new IllegalArgumentException("Name pattern should be specified if 'Entity by name pattern' source is selected."); } EntitiesByNameAndTypeLoader.checkEntityType(EntityType.valueOf(conf.getEntityType())); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java index 6449f832cd..76d42c2ef0 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java @@ -24,13 +24,12 @@ import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; import java.util.Collections; +import static org.thingsboard.rule.engine.transform.OriginatorSource.CUSTOMER; + @Data public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration { - private static final String CUSTOMER_SOURCE = "CUSTOMER"; - - private String originatorSource; - + private OriginatorSource originatorSource; private RelationsQuery relationsQuery; private String entityType; private String entityNamePattern; @@ -38,7 +37,7 @@ public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Originator source should be specified."); + } - ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - ArgumentCaptor originatorCaptor = ArgumentCaptor.forClass(EntityId.class); - verify(ctx).transformMsgOriginator(msgCaptor.capture(), originatorCaptor.capture()); + @Test + public void givenRelatedSourceAndRelatedQueryIsNull_whenInit_thenThrowsException() { + config.setOriginatorSource(RELATED); + config.setRelationsQuery(null); - assertEquals(customerId, originatorCaptor.getValue()); + assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Relations query should be specified if 'Related entity' source is selected."); } @Test - public void newChainCanBeStarted() { - AssetId assetId = new AssetId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - Asset asset = new Asset(); - asset.setCustomerId(customerId); + public void givenEntitySourceAndEntityTypeIsNull_whenInit_thenThrowsException() { + config.setOriginatorSource(ENTITY); + config.setEntityType(null); - RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Entity type should be specified if 'Entity by name pattern' source is selected."); + } - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, assetId, TbMsgMetaData.EMPTY, TbMsgDataType.JSON,TbMsg.EMPTY_JSON_OBJECT, ruleChainId, ruleNodeId); + @ParameterizedTest + @NullAndEmptySource + public void givenEntitySourceAndEntityNamePatternIsEmpty_whenInit_thenThrowsException(String entityName) { + config.setOriginatorSource(ENTITY); + config.setEntityType(EntityType.DEVICE.name()); + config.setEntityNamePattern(entityName); - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); + assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Name pattern should be specified if 'Entity by name pattern' source is selected."); + } - node.onMsg(ctx, msg); - ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - ArgumentCaptor originatorCaptor = ArgumentCaptor.forClass(EntityId.class); - verify(ctx).transformMsgOriginator(msgCaptor.capture(), originatorCaptor.capture()); + @Test + public void givenEntitySourceAndUnexpectedEntityType_whenInit_thenThrowsException() { + config.setOriginatorSource(ENTITY); + config.setEntityType(EntityType.TENANT.name()); + config.setEntityNamePattern("tenant-A"); + + assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Unexpected entity type TENANT"); + } - assertEquals(customerId, originatorCaptor.getValue()); + @Test + public void givenOriginatorSourceIsCustomer_whenOnMsg_thenTellSuccess() throws TbNodeException { + Device device = new Device(DEVICE_ID); + device.setCustomerId(CUSTOMER_ID); + + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, CUSTOMER_ID); + + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getDeviceService()).willReturn(deviceServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(deviceServiceMock.findDeviceById(any(TenantId.class), any(DeviceId.class))).willReturn(device); + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); + + then(deviceServiceMock).should().findDeviceById(TENANT_ID, DEVICE_ID); + then(ctxMock).should().transformMsgOriginator(msg, CUSTOMER_ID); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().tellSuccess(actualMsg.capture()); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg); } @Test - public void exceptionThrownIfCannotFindNewOriginator() { - AssetId assetId = new AssetId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - Asset asset = new Asset(); - asset.setCustomerId(customerId); + public void givenOriginatorSourceIsTenant_whenOnMsg_thenTellSuccess() throws TbNodeException { + config.setOriginatorSource(TENANT); - RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, ASSET_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, TENANT_ID); - TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, assetId, TbMsgMetaData.EMPTY, TbMsgDataType.JSON,TbMsg.EMPTY_JSON_OBJECT, ruleChainId, ruleNodeId); + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(null)); + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); - ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(NoSuchElementException.class); + then(ctxMock).should().transformMsgOriginator(msg, TENANT_ID); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().tellSuccess(actualMsg.capture()); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg); + } - node.onMsg(ctx, msg); - verify(ctx).tellFailure(same(msg), exceptionCaptor.capture()); + @Test + public void givenOriginatorSourceIsRelatedAndNewOriginatorIsNull_whenOnMsg_thenTellFailure() throws TbNodeException { + config.setOriginatorSource(RELATED); + + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, ASSET_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getRelationService()).willReturn(relationServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(relationServiceMock.findByQuery(any(TenantId.class), any(EntityRelationsQuery.class))).willReturn(Futures.immediateFuture(Collections.emptyList())); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); + + var query = new EntityRelationsQuery(); + var relationsQuery = config.getRelationsQuery(); + var parameters = new RelationsSearchParameters( + ASSET_ID, + relationsQuery.getDirection(), + relationsQuery.getMaxLevel(), + relationsQuery.isFetchLastLevelOnly() + ); + query.setParameters(parameters); + query.setFilters(relationsQuery.getFilters()); + then(relationServiceMock).should().findByQuery(TENANT_ID, query); + ArgumentCaptor throwable = ArgumentCaptor.forClass(Throwable.class); + then(ctxMock).should().tellFailure(eq(msg), throwable.capture()); + assertThat(throwable.getValue()).isInstanceOf(NoSuchElementException.class).hasMessage("Failed to find new originator!"); + } - assertEquals("Failed to find new originator!", exceptionCaptor.getValue().getMessage()); + @Test + public void givenOriginatorSourceIsAlarmOriginator_whenOnMsg_thenTellSuccess() throws TbNodeException { + config.setOriginatorSource(ALARM_ORIGINATOR); + + AlarmId alarmId = new AlarmId(UUID.fromString("6b43f694-cb5f-4199-9023-e9e40eeb82dd")); + Alarm alarm = new Alarm(alarmId); + alarm.setOriginator(DEVICE_ID); + + TbMsg msg = TbMsg.newMsg(TbMsgType.ALARM, alarmId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, DEVICE_ID); + + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getAlarmService()).willReturn(alarmServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(alarmServiceMock.findAlarmByIdAsync(any(TenantId.class), any(AlarmId.class))).willReturn(Futures.immediateFuture(alarm)); + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); + + then(alarmServiceMock).should().findAlarmByIdAsync(TENANT_ID, alarmId); + then(ctxMock).should().transformMsgOriginator(msg, DEVICE_ID); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().tellSuccess(actualMsg.capture()); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg); } - public void init() throws TbNodeException { - TbChangeOriginatorNodeConfiguration config = new TbChangeOriginatorNodeConfiguration(); - config.setOriginatorSource(CUSTOMER_SOURCE); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + @ParameterizedTest + @MethodSource + public void givenOriginatorSourceIsEntity_whenOnMsg_thenTellSuccess(String entityNamePattern, TbMsgMetaData metaData, String data) throws TbNodeException { + config.setOriginatorSource(ENTITY); + config.setEntityType(EntityType.ASSET.name()); + config.setEntityNamePattern(entityNamePattern); + + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, data); + TbMsg expectedMsg = TbMsg.transformMsgOriginator(msg, ASSET_ID); + + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getAssetService()).willReturn(assetServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(new Asset(ASSET_ID)); + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); + + String expectedEntityName = TbNodeUtils.processPattern(entityNamePattern, msg); + then(assetServiceMock).should().findAssetByTenantIdAndName(TENANT_ID, expectedEntityName); + then(ctxMock).should().transformMsgOriginator(msg, ASSET_ID); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().tellSuccess(actualMsg.capture()); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg); + } - when(ctx.getDbCallbackExecutor()).thenReturn(dbExecutor); + private static Stream givenOriginatorSourceIsEntity_whenOnMsg_thenTellSuccess() { + return Stream.of( + Arguments.of("test-asset", TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT), + Arguments.of("${md-name-pattern}", new TbMsgMetaData(Map.of("md-name-pattern", "md-test-asset")), TbMsg.EMPTY_JSON_OBJECT), + Arguments.of("${msg-name-pattern}", TbMsgMetaData.EMPTY, "{\"msg-name-pattern\":\"msg-test-asset\"}") + ); + } - node = new TbChangeOriginatorNode(); - node.init(null, nodeConfiguration); + @Test + public void givenOriginatorSourceIsEntityAndEntityCouldNotFound_whenOnMsg_thenTellFailure() throws TbNodeException { + config.setOriginatorSource(ENTITY); + config.setEntityType(EntityType.ASSET.name()); + config.setEntityNamePattern("${md-name-pattern}"); + + TbMsgMetaData metaData = new TbMsgMetaData(); + metaData.putValue("md-name-pattern", "test-asset"); + TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metaData, TbMsg.EMPTY_JSON_OBJECT); + + given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); + given(ctxMock.getAssetService()).willReturn(assetServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); + given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(null); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + node.onMsg(ctxMock, msg); + + ArgumentCaptor throwable = ArgumentCaptor.forClass(Throwable.class); + then(ctxMock).should().tellFailure(eq(msg), throwable.capture()); + assertThat(throwable.getValue()).isInstanceOf(IllegalStateException.class).hasMessage("Failed to find asset with name 'test-asset'!"); } + }