Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard

pull/11466/head
Igor Kulikov 2 years ago
parent
commit
1e1e725f6a
  1. 17
      application/src/main/data/json/system/widget_bundles/scada_water_system_symbols.json
  2. 6
      application/src/main/java/org/thingsboard/server/controller/AdminController.java
  3. 4
      application/src/main/resources/thingsboard.yml
  4. 4
      common/cache/src/main/java/org/thingsboard/server/cache/TbTransactionalCache.java
  5. 6
      common/data/src/main/java/org/thingsboard/server/common/data/objects/AttributesEntityView.java
  6. 2
      common/data/src/main/java/org/thingsboard/server/common/data/objects/TelemetryEntityView.java
  7. 3
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  8. 1
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsTest.java
  9. 3
      msa/tb-node/docker/Dockerfile
  10. 24
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/OriginatorSource.java
  11. 49
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java
  12. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java
  13. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesByNameAndTypeLoader.java
  14. 77
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbAckNodeTest.java
  15. 90
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbCheckpointNodeTest.java
  16. 6
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNodeTest.java
  17. 83
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbRuleChainOutputNodeTest.java
  18. 316
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java

17
application/src/main/data/json/system/widget_bundles/scada_water_system_symbols.json

@ -47,6 +47,21 @@
"vertical_wheel_valve", "vertical_wheel_valve",
"horizontal_ball_valve", "horizontal_ball_valve",
"vertical_ball_valve", "vertical_ball_valve",
"vertical_tank" "vertical_tank",
"stand_vertical_tank",
"cylindrical_tank",
"stand_cylindrical_tank",
"vertical_short_tank",
"stand_vertical_short_tank",
"large_cylindrical_tank",
"large_stand_cylindrical_tank",
"large_vertical_tank",
"large_stand_vertical_tank",
"horizontal_tank",
"stand_horizontal_tank",
"spherical_tank",
"small_spherical_tank",
"elevated_tank",
"pool"
] ]
} }

6
application/src/main/java/org/thingsboard/server/controller/AdminController.java

@ -137,7 +137,7 @@ public class AdminController extends BaseController {
return adminSettings; return adminSettings;
} }
@ApiOperation(value = "Get the Administration Settings object using key (getAdminSettings)", @ApiOperation(value = "Creates or Updates the Administration Settings (saveAdminSettings)",
notes = "Creates or Updates the Administration Settings. Platform generates random Administration Settings Id during settings creation. " + notes = "Creates or Updates the Administration Settings. Platform generates random Administration Settings Id during settings creation. " +
"The Administration Settings Id will be present in the response. Specify the Administration Settings Id when you would like to update the Administration Settings. " + "The Administration Settings Id will be present in the response. Specify the Administration Settings Id when you would like to update the Administration Settings. " +
"Referencing non-existing Administration Settings Id will cause an error." + SYSTEM_AUTHORITY_PARAGRAPH) "Referencing non-existing Administration Settings Id will cause an error." + SYSTEM_AUTHORITY_PARAGRAPH)
@ -160,7 +160,7 @@ public class AdminController extends BaseController {
return adminSettings; return adminSettings;
} }
@ApiOperation(value = "Get the Security Settings object", @ApiOperation(value = "Get the Security Settings object (getSecuritySettings)",
notes = "Get the Security Settings object that contains password policy, etc." + SYSTEM_AUTHORITY_PARAGRAPH) notes = "Get the Security Settings object that contains password policy, etc." + SYSTEM_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAuthority('SYS_ADMIN')") @PreAuthorize("hasAuthority('SYS_ADMIN')")
@RequestMapping(value = "/securitySettings", method = RequestMethod.GET) @RequestMapping(value = "/securitySettings", method = RequestMethod.GET)
@ -237,7 +237,7 @@ public class AdminController extends BaseController {
} }
} }
@ApiOperation(value = "Send test sms (sendTestMail)", @ApiOperation(value = "Send test sms (sendTestSms)",
notes = "Attempts to send test sms to the System Administrator User using SMS Settings and phone number provided as a parameters of the request. " notes = "Attempts to send test sms to the System Administrator User using SMS Settings and phone number provided as a parameters of the request. "
+ SYSTEM_AUTHORITY_PARAGRAPH) + SYSTEM_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAuthority('SYS_ADMIN')") @PreAuthorize("hasAuthority('SYS_ADMIN')")

4
application/src/main/resources/thingsboard.yml

@ -20,8 +20,8 @@ server:
address: "${HTTP_BIND_ADDRESS:0.0.0.0}" address: "${HTTP_BIND_ADDRESS:0.0.0.0}"
# Server bind port # Server bind port
port: "${HTTP_BIND_PORT:8080}" port: "${HTTP_BIND_PORT:8080}"
# Server forward headers strategy # Server forward headers strategy. Required for SWAGGER UI when reverse proxy is used
forward_headers_strategy: "${HTTP_FORWARD_HEADERS_STRATEGY:NONE}" forward_headers_strategy: "${HTTP_FORWARD_HEADERS_STRATEGY:framework}"
# Server SSL configuration # Server SSL configuration
ssl: ssl:
# Enable/disable SSL support # Enable/disable SSL support

4
common/cache/src/main/java/org/thingsboard/server/cache/TbTransactionalCache.java

@ -53,7 +53,7 @@ public interface TbTransactionalCache<K extends Serializable, V extends Serializ
if (putToCache) { if (putToCache) {
return getAndPutInTransaction(key, dbCall, cacheNullValue); return getAndPutInTransaction(key, dbCall, cacheNullValue);
} else { } else {
TbCacheValueWrapper<V> cacheValueWrapper = get(key); TbCacheValueWrapper<V> cacheValueWrapper = get(key, true);
if (cacheValueWrapper != null) { if (cacheValueWrapper != null) {
return cacheValueWrapper.get(); return cacheValueWrapper.get();
} }
@ -92,7 +92,7 @@ public interface TbTransactionalCache<K extends Serializable, V extends Serializ
if (putToCache) { if (putToCache) {
return getAndPutInTransaction(key, dbCall, cacheValueToResult, dbValueToCacheValue, cacheNullValue); return getAndPutInTransaction(key, dbCall, cacheValueToResult, dbValueToCacheValue, cacheNullValue);
} else { } else {
TbCacheValueWrapper<V> cacheValueWrapper = get(key); TbCacheValueWrapper<V> cacheValueWrapper = get(key, true);
if (cacheValueWrapper != null) { if (cacheValueWrapper != null) {
var cacheValue = cacheValueWrapper.get(); var cacheValue = cacheValueWrapper.get();
return cacheValue == null ? null : cacheValueToResult.apply(cacheValue); return cacheValue == null ? null : cacheValueToResult.apply(cacheValue);

6
common/data/src/main/java/org/thingsboard/server/common/data/objects/AttributesEntityView.java

@ -31,11 +31,11 @@ import java.util.List;
@NoArgsConstructor @NoArgsConstructor
public class AttributesEntityView implements Serializable { public class AttributesEntityView implements Serializable {
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of client-side attribute keys to expose", example = "currentConfiguration") @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of client-side attribute keys to expose", example = "[\"currentConfiguration\"]")
private List<String> cs = new ArrayList<>(); private List<String> cs = new ArrayList<>();
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of server-side attribute keys to expose", example = "model") @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of server-side attribute keys to expose", example = "[\"model\"]")
private List<String> ss = new ArrayList<>(); private List<String> ss = new ArrayList<>();
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of shared attribute keys to expose", example = "targetConfiguration") @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of shared attribute keys to expose", example = "[\"targetConfiguration\"]")
private List<String> sh = new ArrayList<>(); private List<String> sh = new ArrayList<>();
public AttributesEntityView(List<String> cs, public AttributesEntityView(List<String> cs,

2
common/data/src/main/java/org/thingsboard/server/common/data/objects/TelemetryEntityView.java

@ -31,7 +31,7 @@ import java.util.List;
@NoArgsConstructor @NoArgsConstructor
public class TelemetryEntityView implements Serializable { public class TelemetryEntityView implements Serializable {
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of time-series data keys to expose", example = "temperature, humidity") @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "List of time-series data keys to expose", example = "[\"temperature\", \"humidity\"]")
private List<String> timeseries; private List<String> timeseries;
@Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "JSON object with attributes to expose") @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "JSON object with attributes to expose")
private AttributesEntityView attributes; private AttributesEntityView attributes;

3
common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java

@ -151,6 +151,7 @@ public class TbKafkaSettings {
Properties props = toProps(); Properties props = toProps();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimeoutMs);
props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, maxPartitionFetchBytes); props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, maxPartitionFetchBytes);
props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, fetchMaxBytes); props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, fetchMaxBytes);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollIntervalMs); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollIntervalMs);
@ -193,8 +194,6 @@ public class TbKafkaSettings {
} }
props.put(CommonClientConfigs.REQUEST_TIMEOUT_MS_CONFIG, requestTimeoutMs); props.put(CommonClientConfigs.REQUEST_TIMEOUT_MS_CONFIG, requestTimeoutMs);
props.put(CommonClientConfigs.SESSION_TIMEOUT_MS_CONFIG, sessionTimeoutMs);
props.putAll(PropertyUtils.getProps(otherInline)); props.putAll(PropertyUtils.getProps(otherInline));
if (other != null) { if (other != null) {

1
common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsTest.java

@ -49,7 +49,6 @@ class TbKafkaSettingsTest {
Properties props = settings.toProps(); Properties props = settings.toProps();
assertThat(props).as("TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS").containsEntry("request.timeout.ms", 30000); assertThat(props).as("TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS").containsEntry("request.timeout.ms", 30000);
assertThat(props).as("TB_QUEUE_KAFKA_SESSION_TIMEOUT_MS").containsEntry("session.timeout.ms", 10000);
//other-inline //other-inline
assertThat(props).as("metrics.recording.level").containsEntry("metrics.recording.level", "INFO"); assertThat(props).as("metrics.recording.level").containsEntry("metrics.recording.level", "INFO");

3
msa/tb-node/docker/Dockerfile

@ -18,9 +18,6 @@ FROM thingsboard/openjdk17:bookworm-slim
COPY start-tb-node.sh ${pkg.name}.deb /tmp/ COPY start-tb-node.sh ${pkg.name}.deb /tmp/
# Required for SWAGGER UI when reverse proxy is used
ENV HTTP_FORWARD_HEADERS_STRATEGY=framework
RUN chmod a+x /tmp/*.sh \ RUN chmod a+x /tmp/*.sh \
&& mv /tmp/start-tb-node.sh /usr/bin && \ && mv /tmp/start-tb-node.sh /usr/bin && \
(yes | dpkg -i /tmp/${pkg.name}.deb) && \ (yes | dpkg -i /tmp/${pkg.name}.deb) && \

24
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
}

49
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; 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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.NoSuchElementException; import java.util.NoSuchElementException;
import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY;
import static org.thingsboard.rule.engine.transform.OriginatorSource.RELATED;
@Slf4j @Slf4j
@RuleNode( @RuleNode(
type = ComponentType.TRANSFORMATION, type = ComponentType.TRANSFORMATION,
@ -59,12 +60,6 @@ import java.util.NoSuchElementException;
) )
public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOriginatorNodeConfiguration> { public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOriginatorNodeConfiguration> {
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 @Override
protected TbChangeOriginatorNodeConfiguration loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { protected TbChangeOriginatorNodeConfiguration loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
var config = TbNodeUtils.convert(configuration, TbChangeOriginatorNodeConfiguration.class); var config = TbNodeUtils.convert(configuration, TbChangeOriginatorNodeConfiguration.class);
@ -85,15 +80,15 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOrig
private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, TbMsg msg) { private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, TbMsg msg) {
switch (config.getOriginatorSource()) { switch (config.getOriginatorSource()) {
case CUSTOMER_SOURCE: case CUSTOMER:
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator()); return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator());
case TENANT_SOURCE: case TENANT:
return Futures.immediateFuture(ctx.getTenantId()); return Futures.immediateFuture(ctx.getTenantId());
case RELATED_SOURCE: case RELATED:
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery()); return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery());
case ALARM_ORIGINATOR_SOURCE: case ALARM_ORIGINATOR:
return EntitiesAlarmOriginatorIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator()); return EntitiesAlarmOriginatorIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator());
case ENTITY_SOURCE: case ENTITY:
EntityType entityType = EntityType.valueOf(config.getEntityType()); EntityType entityType = EntityType.valueOf(config.getEntityType());
String entityName = TbNodeUtils.processPattern(config.getEntityNamePattern(), msg); String entityName = TbNodeUtils.processPattern(config.getEntityNamePattern(), msg);
try { try {
@ -108,28 +103,22 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode<TbChangeOrig
} }
private void validateConfig(TbChangeOriginatorNodeConfiguration conf) { private void validateConfig(TbChangeOriginatorNodeConfiguration conf) {
HashSet<String> knownSources = Sets.newHashSet(CUSTOMER_SOURCE, TENANT_SOURCE, RELATED_SOURCE, ALARM_ORIGINATOR_SOURCE, ENTITY_SOURCE); if (conf.getOriginatorSource() == null) {
if (!knownSources.contains(conf.getOriginatorSource())) { log.debug("Originator source should be specified.");
log.error("Unsupported source [{}] for TbChangeOriginatorNode", conf.getOriginatorSource()); throw new IllegalArgumentException("Originator source should be specified.");
throw new IllegalArgumentException("Unsupported source TbChangeOriginatorNode" + conf.getOriginatorSource());
} }
if (conf.getOriginatorSource().equals(RELATED) && conf.getRelationsQuery() == null) {
if (conf.getOriginatorSource().equals(RELATED_SOURCE)) { log.debug("Relations query should be specified if 'Related entity' source is selected.");
if (conf.getRelationsQuery() == null) { throw new IllegalArgumentException("Relations query should be specified if 'Related entity' source is selected.");
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(ENTITY)) {
if (conf.getOriginatorSource().equals(ENTITY_SOURCE)) {
if (conf.getEntityType() == null) { if (conf.getEntityType() == null) {
log.error("Entity type not specified for [{}]", ENTITY_SOURCE); log.debug("Entity type should be specified if '{}' source is selected.", ENTITY);
throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE); throw new IllegalArgumentException("Entity type should be specified if 'Entity by name pattern' source is selected.");
} }
if (StringUtils.isEmpty(conf.getEntityNamePattern())) { if (StringUtils.isEmpty(conf.getEntityNamePattern())) {
log.error("EntityNamePattern not specified for type [{}]", conf.getEntityType()); log.debug("Name pattern should be specified if '{}' source is selected.", ENTITY);
throw new IllegalArgumentException("Wrong config for [{}] in TbChangeOriginatorNode!" + ENTITY_SOURCE); throw new IllegalArgumentException("Name pattern should be specified if 'Entity by name pattern' source is selected.");
} }
EntitiesByNameAndTypeLoader.checkEntityType(EntityType.valueOf(conf.getEntityType())); EntitiesByNameAndTypeLoader.checkEntityType(EntityType.valueOf(conf.getEntityType()));
} }

9
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 java.util.Collections;
import static org.thingsboard.rule.engine.transform.OriginatorSource.CUSTOMER;
@Data @Data
public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration<TbChangeOriginatorNodeConfiguration> { public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration<TbChangeOriginatorNodeConfiguration> {
private static final String CUSTOMER_SOURCE = "CUSTOMER"; private OriginatorSource originatorSource;
private String originatorSource;
private RelationsQuery relationsQuery; private RelationsQuery relationsQuery;
private String entityType; private String entityType;
private String entityNamePattern; private String entityNamePattern;
@ -38,7 +37,7 @@ public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration<Tb
@Override @Override
public TbChangeOriginatorNodeConfiguration defaultConfiguration() { public TbChangeOriginatorNodeConfiguration defaultConfiguration() {
TbChangeOriginatorNodeConfiguration configuration = new TbChangeOriginatorNodeConfiguration(); TbChangeOriginatorNodeConfiguration configuration = new TbChangeOriginatorNodeConfiguration();
configuration.setOriginatorSource(CUSTOMER_SOURCE); configuration.setOriginatorSource(CUSTOMER);
RelationsQuery relationsQuery = new RelationsQuery(); RelationsQuery relationsQuery = new RelationsQuery();
relationsQuery.setDirection(EntitySearchDirection.FROM); relationsQuery.setDirection(EntitySearchDirection.FROM);

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

@ -53,7 +53,7 @@ public class EntitiesByNameAndTypeLoader {
throw new IllegalStateException("Unexpected entity type " + entityType.name()); throw new IllegalStateException("Unexpected entity type " + entityType.name());
} }
if (targetEntity == null) { if (targetEntity == null) {
throw new IllegalStateException("Failed to found " + entityType.name() + " entity by name: '" + entityName + "'!"); throw new IllegalStateException("Failed to find " + entityType.getNormalName().toLowerCase() + " with name '" + entityName + "'!");
} }
return targetEntity.getId(); return targetEntity.getId();
} }

77
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbAckNodeTest.java

@ -0,0 +1,77 @@
/**
* 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.flow;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
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.id.DeviceId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.mockito.BDDMockito.then;
@ExtendWith(MockitoExtension.class)
public class TbAckNodeTest {
private TbAckNode node;
private EmptyNodeConfiguration config;
private TbNodeConfiguration nodeConfiguration;
@Mock
private TbContext ctxMock;
@BeforeEach
public void setUp() {
node = new TbAckNode();
config = new EmptyNodeConfiguration().defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
}
@Test
public void verifyDefaultConfig() {
assertThat(config.getVersion()).isEqualTo(0);
}
@Test
public void givenDefaultConfig_whenInit_thenOk() {
assertThatNoException().isThrownBy(() -> node.init(ctxMock, nodeConfiguration));
}
@Test
public void givenMsg_whenOnMsg_thenAckAndTellSuccess() throws TbNodeException {
node.init(ctxMock, nodeConfiguration);
DeviceId deviceId = new DeviceId(UUID.fromString("5770153d-6ca2-4447-8a54-5d8a4538e052"));
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, deviceId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
node.onMsg(ctxMock, msg);
then(ctxMock).should().ack(msg);
then(ctxMock).should().tellSuccess(msg);
}
}

90
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbCheckpointNodeTest.java

@ -16,17 +16,103 @@
package org.thingsboard.rule.engine.flow; package org.thingsboard.rule.engine.flow;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest; import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.UUID;
import java.util.function.Consumer;
import java.util.stream.Stream; import java.util.stream.Stream;
import static org.mockito.Mockito.spy; import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.spy;
import static org.mockito.BDDMockito.then;
@Slf4j @Slf4j
@ExtendWith(MockitoExtension.class)
public class TbCheckpointNodeTest extends AbstractRuleNodeUpgradeTest { public class TbCheckpointNodeTest extends AbstractRuleNodeUpgradeTest {
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("37840655-b7dc-4f49-8da3-9429159e0970"));
private TbCheckpointNode node;
private EmptyNodeConfiguration config;
private TbNodeConfiguration nodeConfiguration;
@Mock
private TbContext ctxMock;
@BeforeEach
public void setUp() {
node = spy(new TbCheckpointNode());
config = new EmptyNodeConfiguration().defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
}
@Test
public void verifyDefaultConfig() {
assertThat(config.getVersion()).isEqualTo(0);
}
@Test
public void givenDefaultConfig_whenInit_thenOk() {
assertThatNoException().isThrownBy(() -> node.init(ctxMock, nodeConfiguration));
}
@ParameterizedTest
@ValueSource(strings = {DataConstants.MAIN_QUEUE_NAME, DataConstants.HP_QUEUE_NAME, DataConstants.SQ_QUEUE_NAME, "Custom queue"})
public void givenQueueName_whenOnMsg_thenTransfersMsgToDefinedQueue(String queueName) throws TbNodeException {
given(ctxMock.getQueueName()).willReturn(queueName);
node.init(ctxMock, nodeConfiguration);
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
node.onMsg(ctxMock, msg);
ArgumentCaptor<Runnable> onSuccess = ArgumentCaptor.forClass(Runnable.class);
then(ctxMock).should().enqueueForTellNext(eq(msg), eq(queueName), eq(TbNodeConnectionType.SUCCESS), onSuccess.capture(), any());
onSuccess.getValue().run();
then(ctxMock).should().ack(msg);
}
@Test
public void givenErrorDuringTransfer_whenOnMsg_thenTellFailure() throws TbNodeException {
given(ctxMock.getQueueName()).willReturn(DataConstants.HP_QUEUE_NAME);
node.init(ctxMock, nodeConfiguration);
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
node.onMsg(ctxMock, msg);
ArgumentCaptor<Consumer<Throwable>> onFailure = ArgumentCaptor.forClass(Consumer.class);
then(ctxMock).should().enqueueForTellNext(eq(msg), eq(DataConstants.HP_QUEUE_NAME), eq(TbNodeConnectionType.SUCCESS), any(), onFailure.capture());
String errorMsg = "Something went wrong.";
onFailure.getValue().accept(new RuntimeException(errorMsg));
ArgumentCaptor<Throwable> throwable = ArgumentCaptor.forClass(Throwable.class);
then(ctxMock).should().tellFailure(eq(msg), throwable.capture());
assertThat(throwable.getValue()).isInstanceOf(RuntimeException.class).hasMessage(errorMsg);
}
// Rule nodes upgrade // Rule nodes upgrade
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
return Stream.of( return Stream.of(
@ -50,6 +136,6 @@ public class TbCheckpointNodeTest extends AbstractRuleNodeUpgradeTest {
@Override @Override
protected TbNode getTestNode() { protected TbNode getTestNode() {
return spy(TbCheckpointNode.class); return node;
} }
} }

6
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNodeTest.java

@ -83,6 +83,12 @@ public class TbRuleChainInputNodeTest extends AbstractRuleNodeUpgradeTest {
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
} }
@Test
public void verifyDefaultConfig() {
assertThat(config.getRuleChainId()).isNull();
assertThat(config.isForwardMsgToDefaultRuleChain()).isFalse();
}
@ParameterizedTest @ParameterizedTest
@MethodSource @MethodSource
public void givenValidConfig_whenInit_thenOk(String ruleChainIdStr, boolean forwardMsgToDefaultRuleChain) throws TbNodeException { public void givenValidConfig_whenInit_thenOk(String ruleChainIdStr, boolean forwardMsgToDefaultRuleChain) throws TbNodeException {

83
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/flow/TbRuleChainOutputNodeTest.java

@ -0,0 +1,83 @@
/**
* 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.flow;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
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.id.DeviceId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.spy;
import static org.mockito.BDDMockito.then;
@ExtendWith(MockitoExtension.class)
public class TbRuleChainOutputNodeTest {
private TbRuleChainOutputNode node;
private EmptyNodeConfiguration config;
private TbNodeConfiguration nodeConfiguration;
@Mock
private TbContext ctxMock;
@BeforeEach
public void setUp() {
node = spy(new TbRuleChainOutputNode());
config = new EmptyNodeConfiguration().defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
}
@Test
public void verifyDefaultConfig() {
assertThat(config.getVersion()).isEqualTo(0);
}
@Test
public void givenDefaultConfig_whenInit_thenOk() {
assertThatNoException().isThrownBy(() -> node.init(ctxMock, nodeConfiguration));
}
@Test
public void givenRuleNodeName_whenOnMsg_thenForwardMsgToTheCallerRuleChainWithRelationTypeMatchesWithRuleNodeName() throws TbNodeException {
RuleNode ruleNode = new RuleNode();
ruleNode.setName("test");
given(ctxMock.getSelf()).willReturn(ruleNode);
node.init(ctxMock, nodeConfiguration);
DeviceId deviceId = new DeviceId(UUID.fromString("f514da88-79b3-46da-9f02-1747c5e84f44"));
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, deviceId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
node.onMsg(ctxMock, msg);
then(ctxMock).should().output(msg, "test");
}
}

316
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java

@ -15,139 +15,319 @@
*/ */
package org.thingsboard.rule.engine.transform; package org.thingsboard.rule.engine.transform;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.NullAndEmptySource;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.TestDbCallbackExecutor; import org.thingsboard.rule.engine.TestDbCallbackExecutor;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.data.RelationsQuery;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset; 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.AssetId;
import org.thingsboard.server.common.data.id.CustomerId; 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.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationsQuery;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter;
import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import org.thingsboard.server.common.msg.TbMsg; 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.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.relation.RelationService;
import java.util.Collections;
import java.util.Map;
import java.util.NoSuchElementException; import java.util.NoSuchElementException;
import java.util.UUID;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.same; import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.verify; import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.when; import static org.thingsboard.rule.engine.transform.OriginatorSource.ALARM_ORIGINATOR;
import static org.thingsboard.rule.engine.transform.OriginatorSource.CUSTOMER;
import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY;
import static org.thingsboard.rule.engine.transform.OriginatorSource.RELATED;
import static org.thingsboard.rule.engine.transform.OriginatorSource.TENANT;
@ExtendWith(MockitoExtension.class) @ExtendWith(MockitoExtension.class)
public class TbChangeOriginatorNodeTest { public class TbChangeOriginatorNodeTest {
private static final String CUSTOMER_SOURCE = "CUSTOMER"; private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("79830b6d-4f93-49bd-9b5b-d31ce51da77b"));
private final CustomerId CUSTOMER_ID = new CustomerId(UUID.fromString("c6b2c94b-5517-4f20-bf8e-ae9407eb8a7a"));
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("990605a4-db46-4ed4-942f-e18200453571"));
private final AssetId ASSET_ID = new AssetId(UUID.fromString("55de3f10-1b55-4950-b711-ed132896b260"));
private final ListeningExecutor dbExecutor = new TestDbCallbackExecutor();
private TbChangeOriginatorNode node; private TbChangeOriginatorNode node;
private TbChangeOriginatorNodeConfiguration config;
@Mock @Mock
private TbContext ctx; private TbContext ctxMock;
@Mock @Mock
private AssetService assetService; private AssetService assetServiceMock;
@Mock
private ListeningExecutor dbExecutor; private DeviceService deviceServiceMock;
@Mock
private RelationService relationServiceMock;
@Mock
private RuleEngineAlarmService alarmServiceMock;
@BeforeEach @BeforeEach
public void before() throws TbNodeException { public void before() throws TbNodeException {
dbExecutor = new TestDbCallbackExecutor(); node = new TbChangeOriginatorNode();
init(); config = new TbChangeOriginatorNodeConfiguration().defaultConfiguration();
} }
@Test @Test
public void originatorCanBeChangedToCustomerId() { public void verifyDefaultConfig() {
AssetId assetId = new AssetId(Uuids.timeBased()); var config = new TbChangeOriginatorNodeConfiguration().defaultConfiguration();
CustomerId customerId = new CustomerId(Uuids.timeBased()); assertThat(config.getOriginatorSource()).isEqualTo(CUSTOMER);
Asset asset = new Asset(); RelationsQuery relationsQuery = new RelationsQuery();
asset.setCustomerId(customerId); relationsQuery.setDirection(EntitySearchDirection.FROM);
relationsQuery.setMaxLevel(1);
RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); RelationEntityTypeFilter relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList());
RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter));
assertThat(config.getRelationsQuery()).isEqualTo(relationsQuery);
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, assetId, TbMsgMetaData.EMPTY, TbMsgDataType.JSON, TbMsg.EMPTY_JSON_OBJECT, ruleChainId, ruleNodeId); assertThat(config.getEntityType()).isNull();
assertThat(config.getEntityNamePattern()).isNull();
}
when(ctx.getAssetService()).thenReturn(assetService); @Test
when(assetService.findAssetByIdAsync(any(),eq( assetId))).thenReturn(Futures.immediateFuture(asset)); public void givenRelatedSourceIsNull_whenInit_thenThrowsException() {
config.setOriginatorSource(null);
node.onMsg(ctx, msg); assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Originator source should be specified.");
}
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); @Test
ArgumentCaptor<EntityId> originatorCaptor = ArgumentCaptor.forClass(EntityId.class); public void givenRelatedSourceAndRelatedQueryIsNull_whenInit_thenThrowsException() {
verify(ctx).transformMsgOriginator(msgCaptor.capture(), originatorCaptor.capture()); 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 @Test
public void newChainCanBeStarted() { public void givenEntitySourceAndEntityTypeIsNull_whenInit_thenThrowsException() {
AssetId assetId = new AssetId(Uuids.timeBased()); config.setOriginatorSource(ENTITY);
CustomerId customerId = new CustomerId(Uuids.timeBased()); config.setEntityType(null);
Asset asset = new Asset();
asset.setCustomerId(customerId);
RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); .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); assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); .isInstanceOf(IllegalArgumentException.class)
.hasMessage("Name pattern should be specified if 'Entity by name pattern' source is selected.");
}
node.onMsg(ctx, msg); @Test
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); public void givenEntitySourceAndUnexpectedEntityType_whenInit_thenThrowsException() {
ArgumentCaptor<EntityId> originatorCaptor = ArgumentCaptor.forClass(EntityId.class); config.setOriginatorSource(ENTITY);
verify(ctx).transformMsgOriginator(msgCaptor.capture(), originatorCaptor.capture()); 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<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
} }
@Test @Test
public void exceptionThrownIfCannotFindNewOriginator() { public void givenOriginatorSourceIsTenant_whenOnMsg_thenTellSuccess() throws TbNodeException {
AssetId assetId = new AssetId(Uuids.timeBased()); config.setOriginatorSource(TENANT);
CustomerId customerId = new CustomerId(Uuids.timeBased());
Asset asset = new Asset();
asset.setCustomerId(customerId);
RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, ASSET_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT);
RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); 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); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(null)); node.onMsg(ctxMock, msg);
ArgumentCaptor<NoSuchElementException> exceptionCaptor = ArgumentCaptor.forClass(NoSuchElementException.class); then(ctxMock).should().transformMsgOriginator(msg, TENANT_ID);
ArgumentCaptor<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
}
node.onMsg(ctx, msg); @Test
verify(ctx).tellFailure(same(msg), exceptionCaptor.capture()); 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> 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<TbMsg> actualMsg = ArgumentCaptor.forClass(TbMsg.class);
then(ctxMock).should().tellSuccess(actualMsg.capture());
assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg);
} }
public void init() throws TbNodeException { @ParameterizedTest
TbChangeOriginatorNodeConfiguration config = new TbChangeOriginatorNodeConfiguration(); @MethodSource
config.setOriginatorSource(CUSTOMER_SOURCE); public void givenOriginatorSourceIsEntity_whenOnMsg_thenTellSuccess(String entityNamePattern, TbMsgMetaData metaData, String data) throws TbNodeException {
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); 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<TbMsg> 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<Arguments> 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(); @Test
node.init(null, nodeConfiguration); 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> 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'!");
} }
} }

Loading…
Cancel
Save