From 4721ed275de2d4569674b2ec6994bcb40f5b516b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 13 Aug 2024 14:08:25 +0300 Subject: [PATCH] added checks for null for fetchMode and OrderBy --- .../rule/engine/metadata/FetchMode.java | 22 +++ .../engine/metadata/TbGetTelemetryNode.java | 64 ++++--- .../TbGetTelemetryNodeConfiguration.java | 13 +- .../metadata/TbGetTelemetryNodeTest.java | 167 ++++++++++++++---- 4 files changed, 195 insertions(+), 71 deletions(-) create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchMode.java diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchMode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchMode.java new file mode 100644 index 0000000000..57fee39a3f --- /dev/null +++ b/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 + +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java index 2b87e80127..5faf0e40e2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java +++ b/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 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 entries, TbMsg msg, List 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 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); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java index c4128a4481..3bb9ff1013 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java +++ b/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 { - 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 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 verifyLimitInQuery) throws TbNodeException { + public void givenFetchModeAndLimit_whenOnMsg_thenVerifyLimitInQuery(FetchMode fetchMode, int limit, Consumer verifyLimitInQuery) throws TbNodeException { // GIVEN config.setFetchMode(fetchMode); config.setLimit(limit); @@ -321,15 +335,15 @@ public class TbGetTelemetryNodeTest { private static Stream givenFetchModeAndLimit_whenOnMsg_thenVerifyLimitInQuery() { return Stream.of( Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL, + FetchMode.ALL, 5, (Consumer) query -> assertThat(query.getLimit()).isEqualTo(5)), Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST, + FetchMode.FIRST, TbGetTelemetryNodeConfiguration.MAX_FETCH_SIZE, (Consumer) query -> assertThat(query.getLimit()).isEqualTo(1)), Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_LAST, + FetchMode.LAST, 10, (Consumer) 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 verifyOrderInQuery) throws TbNodeException { + public void givenFetchModeAndOrder_whenOnMsg_thenVerifyOrderInQuery(FetchMode fetchMode, Direction orderBy, Consumer verifyOrderInQuery) throws TbNodeException { // GIVEN config.setFetchMode(fetchMode); config.setOrderBy(orderBy); @@ -361,20 +375,16 @@ public class TbGetTelemetryNodeTest { private static Stream givenFetchModeAndOrder_whenOnMsg_thenVerifyOrderInQuery() { return Stream.of( Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL, - "", - (Consumer) query -> assertThat(query.getOrder()).isEqualTo("ASC")), - Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL, - "DESC", + FetchMode.ALL, + Direction.DESC, (Consumer) query -> assertThat(query.getOrder()).isEqualTo("DESC")), Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST, - "ASC", + FetchMode.FIRST, + Direction.ASC, (Consumer) query -> assertThat(query.getOrder()).isEqualTo("ASC")), Arguments.of( - TbGetTelemetryNodeConfiguration.FETCH_MODE_LAST, - "ASC", + FetchMode.LAST, + Direction.ASC, (Consumer) 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 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; + } }