Browse Source

Merge pull request #7268 from YuriyLytvynchuk/bug/node_originator_telemetry_bugfix

[3.4.2] Bugfix: node 'originator telemetry'
pull/7326/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
884307b513
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 81
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  2. 17
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeTest.java

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

@ -98,7 +98,7 @@ public class TbGetTelemetryNode implements TbNode {
}
Aggregation parseAggregationConfig(String aggName) {
if (StringUtils.isEmpty(aggName)) {
if (StringUtils.isEmpty(aggName) || !fetchMode.equals(FETCH_MODE_ALL)) {
return Aggregation.NONE;
}
return Aggregation.valueOf(aggName);
@ -110,11 +110,9 @@ public class TbGetTelemetryNode implements TbNode {
ctx.tellFailure(msg, new IllegalStateException("Telemetry is not selected!"));
} else {
try {
if (config.isUseMetadataIntervalPatterns()) {
checkMetadataKeyPatterns(msg);
}
Interval interval = getInterval(msg);
List<String> keys = TbNodeUtils.processPatterns(tsKeyNames, msg);
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg, keys));
ListenableFuture<List<TsKvEntry>> 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<ReadTsKvQuery> buildQueries(TbMsg msg, List<String> keys) {
final Interval interval = getInterval(msg);
private List<ReadTsKvQuery> buildQueries(Interval interval, List<String> keys) {
final long aggIntervalStep = Aggregation.NONE.equals(aggregation) ? 1 :
// exact how it validates on BaseTimeseriesService.validate()
// see CassandraBaseTimeseriesDao.findAllAsync()
@ -210,71 +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 checkMetadataKeyPatterns(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 (getMetadataValue(msg, startIntervalPattern) == null && getMetadataValue(msg, endIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata values: '" +
replaceRegex(startIntervalPattern) + "' and '" +
replaceRegex(endIntervalPattern) + "' are undefined");
} else {
if (getMetadataValue(msg, startIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(startIntervalPattern) + "' is undefined");
}
if (getMetadataValue(msg, endIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata 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 metadata values: '" +
replaceRegex(startIntervalPattern) + "' and '" +
replaceRegex(endIntervalPattern) + "' have invalid format");
} else {
if (getInterval(msg).getStartTs() == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(startIntervalPattern) + "' has invalid format");
}
if (getInterval(msg).getEndTs() == null) {
throw new IllegalArgumentException("Message metadata 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 getMetadataValue(TbMsg msg, String pattern) {
return msg.getMetaData().getValue(replaceRegex(pattern));
private String getValuePattern(TbMsg msg, String pattern) {
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) {

17
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());
}

Loading…
Cancel
Save