Browse Source

added checks for null for fetchMode and OrderBy

pull/11087/head
IrynaMatveieva 2 years ago
parent
commit
4721ed275d
  1. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchMode.java
  2. 64
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  3. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java
  4. 167
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

22
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchMode.java

@ -0,0 +1,22 @@
/**
* 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.metadata;
public enum FetchMode {
FIRST, ALL, LAST
}

64
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.ListenableFuture;
@ -35,7 +36,9 @@ import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.SortOrder.Direction;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
@ -50,6 +53,7 @@ import java.util.stream.Collectors;
@RuleNode(type = ComponentType.ENRICHMENT,
name = "originator telemetry",
configClazz = TbGetTelemetryNodeConfiguration.class,
version = 1,
nodeDescription = "Adds message originator telemetry for selected time range into message metadata",
nodeDetails = "Useful when you need to get telemetry data set from the message originator for a specific time range " +
"instead of fetching just the latest telemetry or if you need to get the closest telemetry to the fetch interval start or end. " +
@ -59,14 +63,11 @@ import java.util.stream.Collectors;
configDirective = "tbEnrichmentNodeGetTelemetryFromDatabase")
public class TbGetTelemetryNode implements TbNode {
private static final String DESC_ORDER = "DESC";
private static final String ASC_ORDER = "ASC";
private TbGetTelemetryNodeConfiguration config;
private List<String> tsKeyNames;
private int limit;
private String fetchMode;
private String orderBy;
private FetchMode fetchMode;
private Direction orderBy;
private Aggregation aggregation;
@Override
@ -76,14 +77,17 @@ public class TbGetTelemetryNode implements TbNode {
if (tsKeyNames.isEmpty()) {
throw new TbNodeException("Telemetry is not selected!", true);
}
limit = config.getFetchMode().equals(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL) ? validateLimit(config.getLimit()) : 1;
if (config.getFetchMode() == null) {
throw new TbNodeException("FetchMode should be specified!", true);
}
limit = FetchMode.ALL.equals(config.getFetchMode()) ? validateLimit(config.getLimit()) : 1;
fetchMode = config.getFetchMode();
orderBy = getOrderBy();
aggregation = parseAggregationConfig(config.getAggregation());
}
Aggregation parseAggregationConfig(String aggName) {
if (StringUtils.isEmpty(aggName) || !fetchMode.equals(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL)) {
if (StringUtils.isEmpty(aggName) || !fetchMode.equals(FetchMode.ALL)) {
return Aggregation.NONE;
}
return Aggregation.valueOf(aggName);
@ -107,35 +111,29 @@ public class TbGetTelemetryNode implements TbNode {
interval.getEndTs() - interval.getStartTs();
return keys.stream()
.map(key -> new BaseReadTsKvQuery(key, interval.getStartTs(), interval.getEndTs(), aggIntervalStep, limit, aggregation, orderBy))
.map(key -> new BaseReadTsKvQuery(key, interval.getStartTs(), interval.getEndTs(), aggIntervalStep, limit, aggregation, orderBy.name()))
.collect(Collectors.toList());
}
private String getOrderBy() throws TbNodeException {
private Direction getOrderBy() throws TbNodeException {
switch (fetchMode) {
case TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL:
return getOrderByFetchAll();
case TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST:
return ASC_ORDER;
case ALL:
if (config.getOrderBy() == null) {
throw new TbNodeException("OrderBy should be specified!", true);
}
return config.getOrderBy();
case FIRST:
return Direction.ASC;
case LAST:
return Direction.DESC;
default:
return DESC_ORDER;
}
}
private String getOrderByFetchAll() throws TbNodeException {
String orderBy = config.getOrderBy();
if (ASC_ORDER.equals(orderBy) || DESC_ORDER.equals(orderBy)) {
return orderBy;
}
if (StringUtils.isBlank(orderBy)) {
return ASC_ORDER;
throw new TbNodeException("FetchMode '" + fetchMode + "' is not supported.", true);
}
throw new TbNodeException("Invalid fetch order selected.", true);
}
private TbMsgMetaData updateMetadata(List<TsKvEntry> entries, TbMsg msg, List<String> keys) {
ObjectNode resultNode = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER);
if (TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL.equals(fetchMode)) {
if (FetchMode.ALL.equals(fetchMode)) {
entries.forEach(entry -> processArray(resultNode, entry));
} else {
entries.forEach(entry -> processSingle(resultNode, entry));
@ -231,4 +229,18 @@ public class TbGetTelemetryNode implements TbNode {
private Long endTs;
}
@Override
public TbPair<Boolean, JsonNode> upgrade(int fromVersion, JsonNode oldConfiguration) throws TbNodeException {
boolean hasChanges = false;
switch (fromVersion) {
case 0 -> {
if (oldConfiguration.has("orderBy") && oldConfiguration.get("orderBy").isNull()) {
((ObjectNode) oldConfiguration).put("orderBy", Direction.ASC.name());
hasChanges = true;
}
}
}
return new TbPair<>(hasChanges, oldConfiguration);
}
}

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java

@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.page.SortOrder.Direction;
import java.util.Collections;
import java.util.List;
@ -29,10 +30,6 @@ import java.util.concurrent.TimeUnit;
@Data
public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetTelemetryNodeConfiguration> {
public static final String FETCH_MODE_FIRST = "FIRST";
public static final String FETCH_MODE_LAST = "LAST";
public static final String FETCH_MODE_ALL = "ALL";
public static final int MAX_FETCH_SIZE = 1000;
private int startInterval;
@ -45,8 +42,8 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
private String startIntervalTimeUnit;
private String endIntervalTimeUnit;
private String fetchMode; //FIRST, LAST, ALL
private String orderBy; //ASC, DESC
private FetchMode fetchMode; //FIRST, LAST, ALL
private Direction orderBy; //ASC, DESC
private String aggregation; //MIN, MAX, AVG, SUM, COUNT, NONE;
private int limit;
@ -56,7 +53,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
public TbGetTelemetryNodeConfiguration defaultConfiguration() {
TbGetTelemetryNodeConfiguration configuration = new TbGetTelemetryNodeConfiguration();
configuration.setLatestTsKeyNames(Collections.emptyList());
configuration.setFetchMode("FIRST");
configuration.setFetchMode(FetchMode.FIRST);
configuration.setStartIntervalTimeUnit(TimeUnit.MINUTES.name());
configuration.setStartInterval(2);
configuration.setEndIntervalTimeUnit(TimeUnit.MINUTES.name());
@ -64,7 +61,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
configuration.setUseMetadataIntervalPatterns(false);
configuration.setStartIntervalPattern("");
configuration.setEndIntervalPattern("");
configuration.setOrderBy("ASC");
configuration.setOrderBy(Direction.ASC);
configuration.setAggregation(Aggregation.NONE.name());
configuration.setLimit(MAX_FETCH_SIZE);
return configuration;

167
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

@ -28,8 +28,10 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest;
import org.thingsboard.rule.engine.TestDbCallbackExecutor;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.id.DeviceId;
@ -42,6 +44,7 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvQuery;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.page.SortOrder.Direction;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
@ -64,7 +67,7 @@ import static org.mockito.BDDMockito.then;
import static org.mockito.BDDMockito.willReturn;
@ExtendWith(MockitoExtension.class)
public class TbGetTelemetryNodeTest {
public class TbGetTelemetryNodeTest extends AbstractRuleNodeUpgradeTest {
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("5738401b-9dba-422b-b656-a62fe7431917"));
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("8a8fd749-b2ec-488b-a6c6-fc66614d8686"));
@ -88,7 +91,7 @@ public class TbGetTelemetryNodeTest {
@Test
public void givenAggregationAsString_whenParseAggregation_thenReturnEnum() throws TbNodeException {
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
//compatibility with old configs without "aggregation" parameter
@ -112,7 +115,7 @@ public class TbGetTelemetryNodeTest {
@Test
public void givenAggregationWhiteSpace_whenParseAggregation_thenException() throws TbNodeException {
// GIVEN
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN-THEN
@ -122,7 +125,7 @@ public class TbGetTelemetryNodeTest {
@Test
public void givenAggregationIncorrect_whenParseAggregation_thenException() throws TbNodeException {
// GIVEN
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// WHEN-THEN
@ -139,8 +142,8 @@ public class TbGetTelemetryNodeTest {
assertThat(config.isUseMetadataIntervalPatterns()).isFalse();
assertThat(config.getStartIntervalTimeUnit()).isEqualTo(TimeUnit.MINUTES.name());
assertThat(config.getEndIntervalTimeUnit()).isEqualTo(TimeUnit.MINUTES.name());
assertThat(config.getFetchMode()).isEqualTo(TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST);
assertThat(config.getOrderBy()).isEqualTo("ASC");
assertThat(config.getFetchMode()).isEqualTo(FetchMode.FIRST);
assertThat(config.getOrderBy()).isEqualTo(Direction.ASC);
assertThat(config.getAggregation()).isEqualTo(Aggregation.NONE.name());
assertThat(config.getLimit()).isEqualTo(1000);
assertThat(config.getLatestTsKeyNames()).isEmpty();
@ -159,33 +162,44 @@ public class TbGetTelemetryNodeTest {
.isEqualTo(true);
}
@ParameterizedTest
@ValueSource(ints = {-1, 0, 1, 1001, 2000})
public void givenFetchModeAllAndLimitIsOutOfRange_whenInit_thenThrowsException(int limit) {
@Test
public void givenFetchModeIsNull_whenInit_thenThrowsException() {
// GIVEN
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setLimit(limit);
config.setFetchMode(null);
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("FetchMode should be specified!")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@Test
public void givenFetchModeAllAndOrderByIsNull_whenInit_thenThrowsException() {
// GIVEN
config.setFetchMode(FetchMode.ALL);
config.setOrderBy(null);
// THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Limit should be in a range from 2 to 1000.")
.hasMessage("OrderBy should be specified!")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@ParameterizedTest
@ValueSource(strings = {".ASC", "ascending", "DESCENDING"})
public void givenFetchModeAllAndInvalidOrderBy_whenInit_thenThrowsException(String orderBy) {
@ValueSource(ints = {-1, 0, 1, 1001, 2000})
public void givenFetchModeAllAndLimitIsOutOfRange_whenInit_thenThrowsException(int limit) {
// GIVEN
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setLimit(2);
config.setOrderBy(orderBy);
config.setFetchMode(FetchMode.ALL);
config.setLimit(limit);
// WHEN-THEN
// THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.hasMessage("Invalid fetch order selected.")
.hasMessage("Limit should be in a range from 2 to 1000.")
.extracting(e -> ((TbNodeException) e).isUnrecoverable())
.isEqualTo(true);
}
@ -269,7 +283,7 @@ public class TbGetTelemetryNodeTest {
// GIVEN
config.setStartInterval(5);
config.setEndInterval(1);
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setFetchMode(FetchMode.ALL);
config.setAggregation(aggregation);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -297,7 +311,7 @@ public class TbGetTelemetryNodeTest {
@ParameterizedTest
@MethodSource
public void givenFetchModeAndLimit_whenOnMsg_thenVerifyLimitInQuery(String fetchMode, int limit, Consumer<ReadTsKvQuery> verifyLimitInQuery) throws TbNodeException {
public void givenFetchModeAndLimit_whenOnMsg_thenVerifyLimitInQuery(FetchMode fetchMode, int limit, Consumer<ReadTsKvQuery> verifyLimitInQuery) throws TbNodeException {
// GIVEN
config.setFetchMode(fetchMode);
config.setLimit(limit);
@ -321,15 +335,15 @@ public class TbGetTelemetryNodeTest {
private static Stream<Arguments> givenFetchModeAndLimit_whenOnMsg_thenVerifyLimitInQuery() {
return Stream.of(
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL,
FetchMode.ALL,
5,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getLimit()).isEqualTo(5)),
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST,
FetchMode.FIRST,
TbGetTelemetryNodeConfiguration.MAX_FETCH_SIZE,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getLimit()).isEqualTo(1)),
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_LAST,
FetchMode.LAST,
10,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getLimit()).isEqualTo(1))
);
@ -337,7 +351,7 @@ public class TbGetTelemetryNodeTest {
@ParameterizedTest
@MethodSource
public void givenFetchModeAndOrder_whenOnMsg_thenVerifyOrderInQuery(String fetchMode, String orderBy, Consumer<ReadTsKvQuery> verifyOrderInQuery) throws TbNodeException {
public void givenFetchModeAndOrder_whenOnMsg_thenVerifyOrderInQuery(FetchMode fetchMode, Direction orderBy, Consumer<ReadTsKvQuery> verifyOrderInQuery) throws TbNodeException {
// GIVEN
config.setFetchMode(fetchMode);
config.setOrderBy(orderBy);
@ -361,20 +375,16 @@ public class TbGetTelemetryNodeTest {
private static Stream<Arguments> givenFetchModeAndOrder_whenOnMsg_thenVerifyOrderInQuery() {
return Stream.of(
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL,
"",
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getOrder()).isEqualTo("ASC")),
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL,
"DESC",
FetchMode.ALL,
Direction.DESC,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getOrder()).isEqualTo("DESC")),
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST,
"ASC",
FetchMode.FIRST,
Direction.ASC,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getOrder()).isEqualTo("ASC")),
Arguments.of(
TbGetTelemetryNodeConfiguration.FETCH_MODE_LAST,
"ASC",
FetchMode.LAST,
Direction.ASC,
(Consumer<ReadTsKvQuery>) query -> assertThat(query.getOrder()).isEqualTo("DESC"))
);
}
@ -403,7 +413,7 @@ public class TbGetTelemetryNodeTest {
public void givenFetchModeAll_whenOnMsg_thenTellSuccessAndVerifyMsg() throws TbNodeException {
// GIVEN
config.setLatestTsKeyNames(List.of("temperature", "humidity"));
config.setFetchMode(TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL);
config.setFetchMode(FetchMode.ALL);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -434,7 +444,7 @@ public class TbGetTelemetryNodeTest {
@ValueSource(strings = {"FIRST", "LAST"})
public void givenFetchMode_whenOnMsg_thenTellSuccessAndVerifyMsg(String fetchMode) throws TbNodeException {
// GIVEN
config.setFetchMode(fetchMode);
config.setFetchMode(FetchMode.valueOf(fetchMode));
config.setLatestTsKeyNames(List.of("temperature", "humidity"));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -467,4 +477,87 @@ public class TbGetTelemetryNodeTest {
given(ctxMock.getDbCallbackExecutor()).willReturn(executor);
}
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
return Stream.of(
// config for version 0 (orderBy id null)
Arguments.of(0,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"fetchMode": "ALL",
"orderBy": null,
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""",
true,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"fetchMode": "ALL",
"orderBy": "ASC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
"""),
// config for version 0 (orderBy is specified)
Arguments.of(0,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"fetchMode": "ALL",
"orderBy": "DESC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""",
false,
"""
{
"latestTsKeyNames": [
],
"aggregation": "NONE",
"fetchMode": "ALL",
"orderBy": "DESC",
"limit": 1000,
"useMetadataIntervalPatterns": false,
"startIntervalPattern": "",
"endIntervalPattern": "",
"startInterval": 2,
"startIntervalTimeUnit": "MINUTES",
"endInterval": 1,
"endIntervalTimeUnit": "MINUTES"
}
""")
);
}
@Override
protected TbNode getTestNode() {
return node;
}
}

Loading…
Cancel
Save