From 369af93e51b6e7cb6480af13ce7c0515629541de Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Tue, 13 Sep 2022 17:09:04 +0300 Subject: [PATCH 1/5] fixbug --- .../rule/engine/metadata/TbGetTelemetryNode.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) 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 320502005b..914eb4cddd 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 @@ -91,7 +91,9 @@ public class TbGetTelemetryNode implements TbNode { orderByFetchAll = ASC_ORDER; } aggregation = parseAggregationConfig(config.getAggregation()); - + if (!fetchMode.equals(FETCH_MODE_ALL)) { + aggregation = Aggregation.NONE; + } mapper = new ObjectMapper(); mapper.configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false); mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true); @@ -236,16 +238,16 @@ public class TbGetTelemetryNode implements TbNode { } private void isUndefined(TbMsg msg, String startIntervalPattern, String endIntervalPattern) { - if (getMetadataValue(msg, startIntervalPattern) == null && getMetadataValue(msg, endIntervalPattern) == null) { + if (getValuePattern(msg, startIntervalPattern) == null && getValuePattern(msg, endIntervalPattern) == null) { throw new IllegalArgumentException("Message metadata values: '" + replaceRegex(startIntervalPattern) + "' and '" + replaceRegex(endIntervalPattern) + "' are undefined"); } else { - if (getMetadataValue(msg, startIntervalPattern) == null) { + if (getValuePattern(msg, startIntervalPattern) == null) { throw new IllegalArgumentException("Message metadata value: '" + replaceRegex(startIntervalPattern) + "' is undefined"); } - if (getMetadataValue(msg, endIntervalPattern) == null) { + if (getValuePattern(msg, endIntervalPattern) == null) { throw new IllegalArgumentException("Message metadata value: '" + replaceRegex(endIntervalPattern) + "' is undefined"); } @@ -269,8 +271,9 @@ public class TbGetTelemetryNode implements TbNode { } } - private String getMetadataValue(TbMsg msg, String pattern) { - return msg.getMetaData().getValue(replaceRegex(pattern)); + private String getValuePattern(TbMsg msg, String pattern) { + String valuePattern = TbNodeUtils.processPattern(pattern, msg); + return valuePattern.equals(pattern) ? null : valuePattern; } private String replaceRegex(String pattern) { From 189c1afb01b69faca09b6a05a8c1859005c1a7a3 Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Fri, 16 Sep 2022 10:12:08 +0300 Subject: [PATCH 2/5] refactor code --- .../engine/metadata/TbGetTelemetryNode.java | 22 +++++++++---------- 1 file changed, 10 insertions(+), 12 deletions(-) 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 914eb4cddd..f95909e813 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 @@ -91,16 +91,14 @@ public class TbGetTelemetryNode implements TbNode { orderByFetchAll = ASC_ORDER; } aggregation = parseAggregationConfig(config.getAggregation()); - if (!fetchMode.equals(FETCH_MODE_ALL)) { - aggregation = Aggregation.NONE; - } + mapper = new ObjectMapper(); mapper.configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false); mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true); } Aggregation parseAggregationConfig(String aggName) { - if (StringUtils.isEmpty(aggName)) { + if (StringUtils.isEmpty(aggName) || !fetchMode.equals(FETCH_MODE_ALL)) { return Aggregation.NONE; } return Aggregation.valueOf(aggName); @@ -113,7 +111,7 @@ public class TbGetTelemetryNode implements TbNode { } else { try { if (config.isUseMetadataIntervalPatterns()) { - checkMetadataKeyPatterns(msg); + checkKeyPatterns(msg); } List keys = TbNodeUtils.processPatterns(tsKeyNames, msg); ListenableFuture> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg, keys)); @@ -232,23 +230,23 @@ public class TbGetTelemetryNode implements TbNode { return NumberUtils.isParsable(TbNodeUtils.processPattern(pattern, msg)); } - private void checkMetadataKeyPatterns(TbMsg msg) { + private void checkKeyPatterns(TbMsg msg) { isUndefined(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern()); isInvalid(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern()); } private void isUndefined(TbMsg msg, String startIntervalPattern, String endIntervalPattern) { if (getValuePattern(msg, startIntervalPattern) == null && getValuePattern(msg, endIntervalPattern) == null) { - throw new IllegalArgumentException("Message metadata values: '" + + throw new IllegalArgumentException("Message values: '" + replaceRegex(startIntervalPattern) + "' and '" + replaceRegex(endIntervalPattern) + "' are undefined"); } else { if (getValuePattern(msg, startIntervalPattern) == null) { - throw new IllegalArgumentException("Message metadata value: '" + + throw new IllegalArgumentException("Message value: '" + replaceRegex(startIntervalPattern) + "' is undefined"); } if (getValuePattern(msg, endIntervalPattern) == null) { - throw new IllegalArgumentException("Message metadata value: '" + + throw new IllegalArgumentException("Message value: '" + replaceRegex(endIntervalPattern) + "' is undefined"); } } @@ -256,16 +254,16 @@ public class TbGetTelemetryNode implements TbNode { private void isInvalid(TbMsg msg, String startIntervalPattern, String endIntervalPattern) { if (getInterval(msg).getStartTs() == null && getInterval(msg).getEndTs() == null) { - throw new IllegalArgumentException("Message metadata values: '" + + throw new IllegalArgumentException("Message values: '" + replaceRegex(startIntervalPattern) + "' and '" + replaceRegex(endIntervalPattern) + "' have invalid format"); } else { if (getInterval(msg).getStartTs() == null) { - throw new IllegalArgumentException("Message metadata value: '" + + throw new IllegalArgumentException("Message value: '" + replaceRegex(startIntervalPattern) + "' has invalid format"); } if (getInterval(msg).getEndTs() == null) { - throw new IllegalArgumentException("Message metadata value: '" + + throw new IllegalArgumentException("Message value: '" + replaceRegex(endIntervalPattern) + "' has invalid format"); } } From d8135f0782b811fb391dee392b8af976a1b0211d Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Fri, 16 Sep 2022 16:11:25 +0300 Subject: [PATCH 3/5] refactor getInterval logic --- .../engine/metadata/TbGetTelemetryNode.java | 86 ++++++------------- 1 file changed, 28 insertions(+), 58 deletions(-) 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 f95909e813..478e49f004 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 @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 The Thingsboard Authors - * + *

* Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * + *

+ * 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. @@ -110,11 +110,9 @@ public class TbGetTelemetryNode implements TbNode { ctx.tellFailure(msg, new IllegalStateException("Telemetry is not selected!")); } else { try { - if (config.isUseMetadataIntervalPatterns()) { - checkKeyPatterns(msg); - } + Interval interval = getInterval(msg); List keys = TbNodeUtils.processPatterns(tsKeyNames, msg); - ListenableFuture> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg, keys)); + ListenableFuture> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(interval, keys)); DonAsynchron.withCallback(list, data -> { process(data, msg, keys); ctx.tellSuccess(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), msg.getData())); @@ -129,8 +127,7 @@ public class TbGetTelemetryNode implements TbNode { public void destroy() { } - private List buildQueries(TbMsg msg, List keys) { - final Interval interval = getInterval(msg); + private List buildQueries(Interval interval, List keys) { final long aggIntervalStep = Aggregation.NONE.equals(aggregation) ? 1 : // exact how it validates on BaseTimeseriesService.validate() // see CassandraBaseTimeseriesDao.findAllAsync() @@ -210,72 +207,45 @@ public class TbGetTelemetryNode implements TbNode { } private Interval getInterval(TbMsg msg) { - Interval interval = new Interval(); if (config.isUseMetadataIntervalPatterns()) { - if (isParsable(msg, config.getStartIntervalPattern())) { - interval.setStartTs(Long.parseLong(TbNodeUtils.processPattern(config.getStartIntervalPattern(), msg))); - } - if (isParsable(msg, config.getEndIntervalPattern())) { - interval.setEndTs(Long.parseLong(TbNodeUtils.processPattern(config.getEndIntervalPattern(), msg))); - } + return getIntervalFromPatterns(msg); } else { + Interval interval = new Interval(); long ts = System.currentTimeMillis(); interval.setStartTs(ts - TimeUnit.valueOf(config.getStartIntervalTimeUnit()).toMillis(config.getStartInterval())); interval.setEndTs(ts - TimeUnit.valueOf(config.getEndIntervalTimeUnit()).toMillis(config.getEndInterval())); + return interval; } - return interval; } - private boolean isParsable(TbMsg msg, String pattern) { - return NumberUtils.isParsable(TbNodeUtils.processPattern(pattern, msg)); - } - - private void checkKeyPatterns(TbMsg msg) { - isUndefined(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern()); - isInvalid(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern()); + private Interval getIntervalFromPatterns(TbMsg msg) { + Interval interval = new Interval(); + interval.setStartTs(checkPattern(msg, config.getStartIntervalPattern())); + interval.setEndTs(checkPattern(msg, config.getEndIntervalPattern())); + return interval; } - private void isUndefined(TbMsg msg, String startIntervalPattern, String endIntervalPattern) { - if (getValuePattern(msg, startIntervalPattern) == null && getValuePattern(msg, endIntervalPattern) == null) { - throw new IllegalArgumentException("Message values: '" + - replaceRegex(startIntervalPattern) + "' and '" + - replaceRegex(endIntervalPattern) + "' are undefined"); - } else { - if (getValuePattern(msg, startIntervalPattern) == null) { - throw new IllegalArgumentException("Message value: '" + - replaceRegex(startIntervalPattern) + "' is undefined"); - } - if (getValuePattern(msg, endIntervalPattern) == null) { - throw new IllegalArgumentException("Message value: '" + - replaceRegex(endIntervalPattern) + "' is undefined"); - } + private long checkPattern(TbMsg msg, String pattern) { + String value = getValuePattern(msg, pattern); + if (value == null) { + throw new IllegalArgumentException("Message value: '" + + replaceRegex(pattern) + "' is undefined"); } - } - - private void isInvalid(TbMsg msg, String startIntervalPattern, String endIntervalPattern) { - if (getInterval(msg).getStartTs() == null && getInterval(msg).getEndTs() == null) { - throw new IllegalArgumentException("Message values: '" + - replaceRegex(startIntervalPattern) + "' and '" + - replaceRegex(endIntervalPattern) + "' have invalid format"); - } else { - if (getInterval(msg).getStartTs() == null) { - throw new IllegalArgumentException("Message value: '" + - replaceRegex(startIntervalPattern) + "' has invalid format"); - } - if (getInterval(msg).getEndTs() == null) { - throw new IllegalArgumentException("Message value: '" + - replaceRegex(endIntervalPattern) + "' has invalid format"); - } + boolean parsable = NumberUtils.isParsable(value); + if (!parsable) { + throw new IllegalArgumentException("Message value: '" + + replaceRegex(pattern) + "' has invalid format"); } + return Long.parseLong(value); } private String getValuePattern(TbMsg msg, String pattern) { - String valuePattern = TbNodeUtils.processPattern(pattern, msg); - return valuePattern.equals(pattern) ? null : valuePattern; + String value = TbNodeUtils.processPattern(pattern, msg); + return value.equals(pattern) ? null : value; } private String replaceRegex(String pattern) { - return pattern.replaceAll("[${}]", ""); + return pattern.replaceAll("[$\\[{}\\]]", ""); } private int validateLimit(int limit) { From 225f572b7aec3a40634b12fdb911280812505eda Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Fri, 16 Sep 2022 16:15:08 +0300 Subject: [PATCH 4/5] license:format --- .../rule/engine/metadata/TbGetTelemetryNode.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 478e49f004..e3f4613d33 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 @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 The Thingsboard Authors - *

+ * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

+ * + * 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. From 14898ef7b3204015883d44e6024d229410239052 Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Wed, 21 Sep 2022 12:53:30 +0300 Subject: [PATCH 5/5] update test --- .../engine/metadata/TbGetTelemetryNodeTest.java | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java index bb086af15a..9c97f6f85b 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java @@ -15,8 +15,12 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.Before; import org.junit.Test; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.TbContext; +import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.server.common.data.kv.Aggregation; import static org.hamcrest.CoreMatchers.is; @@ -24,14 +28,25 @@ import static org.hamcrest.MatcherAssert.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.BDDMockito.willCallRealMethod; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; public class TbGetTelemetryNodeTest { TbGetTelemetryNode node; + TbGetTelemetryNodeConfiguration config; + TbNodeConfiguration nodeConfiguration; + TbContext ctx; @Before public void setUp() throws Exception { - node = mock(TbGetTelemetryNode.class); + ctx = mock(TbContext.class); + node = spy(new TbGetTelemetryNode()); + config = new TbGetTelemetryNodeConfiguration(); + config.setFetchMode("ALL"); + ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; + nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + node.init(ctx, nodeConfiguration); + willCallRealMethod().given(node).parseAggregationConfig(any()); }