Browse Source

added more tests

pull/14141/head
IrynaMatveieva 12 months ago
parent
commit
a693cabc05
  1. 46
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 6
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 3
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  4. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java
  5. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  7. 10
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java
  8. 9
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java
  9. 8
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggArgumentEntry.java
  10. 17
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleEntityArgumentEntry.java
  11. 103
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/LatestValuesAggregationCalculatedFieldState.java
  12. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/AvgAggEntry.java
  13. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/new_agg.json
  14. 4
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java
  15. 10
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java
  16. 344
      application/src/test/java/org/thingsboard/server/cf/LatestValuesAggregationCalculatedFieldTest.java
  17. 1
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/LatestValuesAggregationCalculatedFieldConfiguration.java

46
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -53,7 +53,7 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleEntityArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.LatestValuesAggregationCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
@ -67,6 +67,7 @@ import java.util.HashSet;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@ -329,7 +330,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} else if (proto.getAttrDataCount() > 0) {
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto));
} else if (proto.getRemovedTsKeysCount() > 0) {
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArgumentsWithFetchedValue(ctx, proto.getRemovedTsKeysList()), toTbMsgId(proto), toTbMsgType(proto));
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArgumentsWithFetchedValue(ctx, msg.getEntityId(), proto.getRemovedTsKeysList()), toTbMsgId(proto), toTbMsgType(proto));
} else if (proto.getRemovedAttrKeysCount() > 0) {
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArgumentsWithDefaultValue(ctx, msg.getEntityId(), proto.getScope(), proto.getRemovedAttrKeysList()), toTbMsgId(proto), toTbMsgType(proto));
} else {
@ -405,7 +406,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
}
private void processRemovedTelemetry(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List<CalculatedFieldId> cfIdList, MultipleTbCallback callback) throws CalculatedFieldException {
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArgumentsWithFetchedValue(ctx, proto.getRemovedTsKeysList()), toTbMsgId(proto), toTbMsgType(proto));
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArgumentsWithFetchedValue(ctx, entityId, proto.getRemovedTsKeysList()), toTbMsgId(proto), toTbMsgType(proto));
}
private void processRemovedAttributes(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List<CalculatedFieldId> cfIdList, MultipleTbCallback callback) throws CalculatedFieldException {
@ -565,7 +566,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
ReferencedEntityKey key = new ReferencedEntityKey(item.getKv().getKey(), ArgumentType.TS_LATEST, null);
String argName = aggArgNames.get(key);
if (argName != null) {
arguments.put(argName, new AggSingleArgumentEntry(originator, item));
arguments.put(argName, new AggSingleEntityArgumentEntry(originator, item));
}
}
}
@ -618,7 +619,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
ReferencedEntityKey key = new ReferencedEntityKey(item.getKey(), ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name()));
String argName = aggArgNames.get(key);
if (argName != null) {
arguments.put(argName, new AggSingleArgumentEntry(entityId, item));
arguments.put(argName, new AggSingleEntityArgumentEntry(entityId, item));
}
}
}
@ -631,14 +632,21 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return Collections.emptyMap();
}
List<String> geofencingArgumentNames = ctx.getLinkedEntityAndCurrentOwnerGeofencingArgumentNames();
return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), geofencingArgumentNames, scope, removedAttrKeys);
List<String> relatedArgumentNames = ctx.getRelatedEntityArgumentNames();
return mapToArgumentsWithDefaultValue(entityId, argNames, ctx.getArguments(), geofencingArgumentNames, relatedArgumentNames, scope, removedAttrKeys);
}
private Map<String, ArgumentEntry> mapToArgumentsWithDefaultValue(CalculatedFieldCtx ctx, AttributeScopeProto scope, List<String> removedAttrKeys) {
return mapToArgumentsWithDefaultValue(ctx.getMainEntityArguments(), ctx.getArguments(), ctx.getMainEntityGeofencingArgumentNames(), scope, removedAttrKeys);
return mapToArgumentsWithDefaultValue(null, ctx.getMainEntityArguments(), ctx.getArguments(), ctx.getMainEntityGeofencingArgumentNames(), new ArrayList<>(), scope, removedAttrKeys);
}
private Map<String, ArgumentEntry> mapToArgumentsWithDefaultValue(Map<ReferencedEntityKey, String> argNames, Map<String, Argument> configArguments, List<String> geofencingArgNames, AttributeScopeProto scope, List<String> removedAttrKeys) {
private Map<String, ArgumentEntry> mapToArgumentsWithDefaultValue(EntityId msgEntityId,
Map<ReferencedEntityKey, String> argNames,
Map<String, Argument> configArguments,
List<String> geofencingArgNames,
List<String> relatedEntityArgNames,
AttributeScopeProto scope,
List<String> removedAttrKeys) {
Map<String, ArgumentEntry> arguments = new HashMap<>();
for (String removedKey : removedAttrKeys) {
ReferencedEntityKey key = new ReferencedEntityKey(removedKey, ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name()));
@ -652,22 +660,36 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
}
Argument argument = configArguments.get(argName);
String defaultValue = (argument != null) ? argument.getDefaultValue() : null;
arguments.put(argName, StringUtils.isNotEmpty(defaultValue)
SingleValueArgumentEntry argumentEntry = StringUtils.isNotEmpty(defaultValue)
? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null)
: new SingleValueArgumentEntry());
: new SingleValueArgumentEntry();
if (relatedEntityArgNames.contains(argName)) {
arguments.put(argName, new AggSingleEntityArgumentEntry(msgEntityId, argumentEntry));
continue;
}
arguments.put(argName, argumentEntry);
}
return arguments;
}
private Map<String, ArgumentEntry> mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, List<String> removedTelemetryKeys) {
private Map<String, ArgumentEntry> mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, EntityId entityId, List<String> removedTelemetryKeys) {
Map<String, Argument> deletedArguments = ctx.getArguments().entrySet().stream()
.filter(entry -> removedTelemetryKeys.contains(entry.getValue().getRefEntityKey().getKey()))
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, entityId, deletedArguments);
fetchedArgs.values().forEach(arg -> arg.setForceResetPrevious(true));
if (CalculatedFieldType.LATEST_VALUES_AGGREGATION.equals(ctx.getCfType())) {
fetchedArgs = fetchedArgs.entrySet().stream()
.collect(Collectors.toMap(
Map.Entry::getKey,
argEntry -> new AggSingleEntityArgumentEntry(entityId, argEntry.getValue())
));
} else {
fetchedArgs.values().forEach(arg -> arg.setForceResetPrevious(true));
}
return fetchedArgs;
}

6
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -516,12 +516,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
.toList();
}
private List<CalculatedFieldCtx> getCfsWithRelationToEntity(EntityId entityId) {
return aggCalculatedFields.values().stream()
.filter(cf -> !findRelationsForCf(entityId, cf).isEmpty())
.toList();
}
private List<CalculatedFieldEntityCtxId> findRelationsForCf(EntityId entityId, CalculatedFieldCtx cf) {
List<CalculatedFieldEntityCtxId> result = new ArrayList<>();
if (cf.getCalculatedField().getConfiguration() instanceof LatestValuesAggregationCalculatedFieldConfiguration configuration) {

3
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -250,7 +250,8 @@ public abstract class AbstractCalculatedFieldProcessingService {
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))), MoreExecutors.directExecutor());
}
public ListenableFuture<ArgumentEntry> fetchAggArgumentEntry(TenantId tenantId, List<EntityId> aggEntities, Argument argument, long startTs) {List<ListenableFuture<Map.Entry<EntityId, ArgumentEntry>>> futures = aggEntities.stream()
public ListenableFuture<ArgumentEntry> fetchAggArgumentEntry(TenantId tenantId, List<EntityId> aggEntities, Argument argument, long startTs) {
List<ListenableFuture<Map.Entry<EntityId, ArgumentEntry>>> futures = aggEntities.stream()
.map(entityId -> {
ListenableFuture<ArgumentEntry> singleAggEntryFut = fetchSingleAggArgumentEntry(tenantId, entityId, argument, startTs);
return Futures.transform(singleAggEntryFut, singleAggEntry -> Map.entry(entityId, singleAggEntry), MoreExecutors.directExecutor());

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java

@ -23,7 +23,7 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleEntityArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
import java.util.List;
@ -39,7 +39,7 @@ import java.util.Map;
@JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING"),
@JsonSubTypes.Type(value = GeofencingArgumentEntry.class, name = "GEOFENCING"),
@JsonSubTypes.Type(value = AggArgumentEntry.class, name = "AGGREGATE_LATEST"),
@JsonSubTypes.Type(value = AggSingleArgumentEntry.class, name = "AGGREGATE_LATEST_SINGLE")
@JsonSubTypes.Type(value = AggSingleEntityArgumentEntry.class, name = "AGGREGATE_LATEST_SINGLE")
})
public interface ArgumentEntry {
@ -75,7 +75,7 @@ public interface ArgumentEntry {
}
static ArgumentEntry createAggSingleArgument(EntityId entityId, KvEntry kvEntry) {
return new AggSingleArgumentEntry(entityId, kvEntry);
return new AggSingleEntityArgumentEntry(entityId, kvEntry);
}
}

11
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java

@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf.ctx.state;
import lombok.Getter;
import lombok.Setter;
import org.thingsboard.script.api.tbel.TbUtils;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -134,4 +135,14 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
this.latestTimestamp = Math.max(this.latestTimestamp, newTs);
}
protected Object formatResult(double result, Integer decimals) {
if (decimals == null) {
return result;
}
if (decimals.equals(0)) {
return TbUtils.toInt(result);
}
return TbUtils.toFixed(result, decimals);
}
}

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -107,6 +107,7 @@ public class CalculatedFieldCtx {
private boolean relationQueryDynamicArguments;
private List<String> mainEntityGeofencingArgumentNames;
private List<String> linkedEntityAndCurrentOwnerGeofencingArgumentNames;
private List<String> relatedEntityArgumentNames;
private long scheduledUpdateIntervalMillis;
@ -126,6 +127,7 @@ public class CalculatedFieldCtx {
this.argNames = new ArrayList<>();
this.mainEntityGeofencingArgumentNames = new ArrayList<>();
this.linkedEntityAndCurrentOwnerGeofencingArgumentNames = new ArrayList<>();
this.relatedEntityArgumentNames = new ArrayList<>();
this.output = calculatedField.getConfiguration().getOutput();
if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argBasedConfig) {
this.arguments.putAll(argBasedConfig.getArguments());
@ -153,6 +155,7 @@ public class CalculatedFieldCtx {
}
}
this.argNames.addAll(arguments.keySet());
this.relatedEntityArgumentNames.addAll(relatedEntityArguments.values());
if (argBasedConfig instanceof ExpressionBasedCalculatedFieldConfiguration expressionBasedConfig) {
this.expression = expressionBasedConfig.getExpression();
this.useLatestTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) argBasedConfig).isUseLatestTs();
@ -174,7 +177,7 @@ public class CalculatedFieldCtx {
}
this.requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation();
if (calculatedField.getConfiguration() instanceof LatestValuesAggregationCalculatedFieldConfiguration aggConfig) {
this.scheduledUpdateIntervalMillis = aggConfig.getDeduplicationIntervalMillis();
this.useLatestTs = aggConfig.isUseLatestTs();
}
this.systemContext = systemContext;
this.tbelInvokeService = systemContext.getTbelInvokeService();
@ -578,6 +581,7 @@ public class CalculatedFieldCtx {
}
if (calculatedField.getConfiguration() instanceof LatestValuesAggregationCalculatedFieldConfiguration thisConfig
&& other.getCalculatedField().getConfiguration() instanceof LatestValuesAggregationCalculatedFieldConfiguration otherConfig
&& thisConfig.getDeduplicationIntervalMillis() != otherConfig.getDeduplicationIntervalMillis()
&& !thisConfig.getMetrics().equals(otherConfig.getMetrics())) {
return true;
}

10
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java

@ -62,16 +62,6 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState {
.build());
}
private Object formatResult(double expressionResult, Integer decimals) {
if (decimals == null) {
return expressionResult;
}
if (decimals.equals(0)) {
return TbUtils.toInt(expressionResult);
}
return TbUtils.toFixed(expressionResult, decimals);
}
private JsonNode createResultJson(boolean useLatestTs, String outputName, Object result) {
ObjectNode valuesNode = JacksonUtil.newObjectNode();
if (result instanceof Double doubleValue) {

9
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java

@ -45,6 +45,15 @@ public class SingleValueArgumentEntry implements ArgumentEntry {
public static final Long DEFAULT_VERSION = -1L;
public SingleValueArgumentEntry(ArgumentEntry entry) {
if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) {
this.ts = singleValueArgumentEntry.ts;
this.kvEntryValue = singleValueArgumentEntry.kvEntryValue;
this.version = singleValueArgumentEntry.version;
this.forceResetPrevious = singleValueArgumentEntry.forceResetPrevious;
}
}
public SingleValueArgumentEntry(TsKvProto entry) {
this.ts = entry.getTs();
if (entry.hasVersion()) {

8
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggArgumentEntry.java

@ -48,11 +48,11 @@ public class AggArgumentEntry implements ArgumentEntry {
if (entry instanceof AggArgumentEntry aggArgumentEntry) {
aggInputs.putAll(aggArgumentEntry.aggInputs);
return true;
} else if (entry instanceof AggSingleArgumentEntry aggSingleArgumentEntry) {
if (aggSingleArgumentEntry.isDeleted()) {
aggInputs.remove(aggSingleArgumentEntry.getEntityId());
} else if (entry instanceof AggSingleEntityArgumentEntry aggSingleEntityArgumentEntry) {
if (aggSingleEntityArgumentEntry.isDeleted()) {
aggInputs.remove(aggSingleEntityArgumentEntry.getEntityId());
} else {
aggInputs.put(aggSingleArgumentEntry.getEntityId(), aggSingleArgumentEntry);
aggInputs.put(aggSingleEntityArgumentEntry.getEntityId(), aggSingleEntityArgumentEntry);
}
return true;
} else {

17
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleArgumentEntry.java → application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleEntityArgumentEntry.java

@ -30,34 +30,39 @@ import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
@Data
@NoArgsConstructor
@AllArgsConstructor
public class AggSingleArgumentEntry extends SingleValueArgumentEntry {
public class AggSingleEntityArgumentEntry extends SingleValueArgumentEntry {
private EntityId entityId;
private boolean deleted;
public AggSingleArgumentEntry(EntityId entityId, TsKvProto entry) {
public AggSingleEntityArgumentEntry(EntityId entityId, ArgumentEntry entry) {
super(entry);
this.entityId = entityId;
}
public AggSingleArgumentEntry(EntityId entityId, AttributeValueProto entry) {
public AggSingleEntityArgumentEntry(EntityId entityId, TsKvProto entry) {
super(entry);
this.entityId = entityId;
}
public AggSingleArgumentEntry(EntityId entityId, KvEntry entry) {
public AggSingleEntityArgumentEntry(EntityId entityId, AttributeValueProto entry) {
super(entry);
this.entityId = entityId;
}
public AggSingleArgumentEntry(EntityId entityId, long ts, BasicKvEntry kvEntryValue, Long version) {
public AggSingleEntityArgumentEntry(EntityId entityId, KvEntry entry) {
super(entry);
this.entityId = entityId;
}
public AggSingleEntityArgumentEntry(EntityId entityId, long ts, BasicKvEntry kvEntryValue, Long version) {
super(ts, kvEntryValue, version);
this.entityId = entityId;
}
@Override
public boolean updateEntry(ArgumentEntry entry) {
if (entry instanceof AggSingleArgumentEntry singleValueEntry) {
if (entry instanceof AggSingleEntityArgumentEntry singleValueEntry) {
if (singleValueEntry.getTs() <= ts) {
return false;
}

103
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/LatestValuesAggregationCalculatedFieldState.java

@ -15,10 +15,12 @@
*/
package org.thingsboard.server.service.cf.ctx.state.aggregation;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.Data;
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.actors.TbActorRef;
@ -43,10 +45,11 @@ import java.util.Map;
import java.util.Map.Entry;
@Slf4j
@Data
@Getter
public class LatestValuesAggregationCalculatedFieldState extends BaseCalculatedFieldState {
private long lastArgsRefreshTs = -1;
@Setter
private long lastMetricsEvalTs = -1;
private long deduplicationInterval = -1;
private Map<String, AggMetric> metrics;
@ -76,8 +79,7 @@ public class LatestValuesAggregationCalculatedFieldState extends BaseCalculatedF
@Override
public void init() {
super.init();
// long scheduledUpdateIntervalMillis = ctx.getScheduledUpdateIntervalMillis();
// ctx.scheduleReevaluation(scheduledUpdateIntervalMillis, actorCtx);
ctx.scheduleReevaluation(deduplicationInterval, actorCtx);
}
@Override
@ -102,41 +104,51 @@ public class LatestValuesAggregationCalculatedFieldState extends BaseCalculatedF
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception {
if (!shouldRecalculate()) {
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.result(null)
.build());
}
Output output = ctx.getOutput();
ObjectNode aggResult = aggregateMetrics(output);
lastMetricsEvalTs = System.currentTimeMillis();
ctx.scheduleReevaluation(deduplicationInterval, actorCtx);
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.type(output.getType())
.scope(output.getScope())
.result(createResultJson(ctx.isUseLatestTs(), aggResult))
.build());
}
private boolean shouldRecalculate() {
boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationInterval;
boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs;
if (intervalPassed && argsUpdatedDuringInterval) {
ObjectNode aggResult = JacksonUtil.newObjectNode();
for (Entry<String, AggMetric> entry : metrics.entrySet()) {
String metricKey = entry.getKey();
AggMetric metric = entry.getValue();
AggEntry aggMetric = AggFunctionFactory.createAggFunction(metric.getFunction());
for (Map<String, ArgumentEntry> entityInputs : inputs.values()) {
if (applyAggregation(metric.getFilter(), entityInputs)) {
Object arg = resolveAggregationInput(metric.getInput(), entityInputs);
if (arg != null) {
aggMetric.update(arg);
}
}
}
return intervalPassed && argsUpdatedDuringInterval;
}
private ObjectNode aggregateMetrics(Output output) throws Exception {
ObjectNode aggResult = JacksonUtil.newObjectNode();
for (Entry<String, AggMetric> entry : metrics.entrySet()) {
String metricKey = entry.getKey();
AggMetric metric = entry.getValue();
aggMetric.result().ifPresent(result -> {
aggResult.set(metricKey, JacksonUtil.valueToTree(result));
});
AggEntry aggMetricEntry = AggFunctionFactory.createAggFunction(metric.getFunction());
aggregateMetric(metric, aggMetricEntry);
aggMetricEntry.result().ifPresent(result -> {
aggResult.set(metricKey, JacksonUtil.valueToTree(formatResult(result, output.getDecimalsByDefault())));
});
}
return aggResult;
}
private void aggregateMetric(AggMetric metric, AggEntry aggEntry) throws Exception {
for (Map<String, ArgumentEntry> entityInputs : inputs.values()) {
if (applyAggregation(metric.getFilter(), entityInputs)) {
Object arg = resolveAggregationInput(metric.getInput(), entityInputs);
if (arg != null) {
aggEntry.update(arg);
}
}
Output output = ctx.getOutput();
lastMetricsEvalTs = System.currentTimeMillis();
ctx.scheduleReevaluation(deduplicationInterval, actorCtx);
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.type(output.getType())
.scope(output.getScope())
.result(aggResult)
.build());
} else {
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.result(null)
.build());
}
}
@ -158,4 +170,25 @@ public class LatestValuesAggregationCalculatedFieldState extends BaseCalculatedF
}
}
private Object formatResult(Object aggregationResult, Integer decimals) {
try {
double result = Double.parseDouble(aggregationResult.toString());
return formatResult(result, decimals);
} catch (Exception e) {
throw new IllegalArgumentException("Aggregation result cannot be parsed: " + aggregationResult, e);
}
}
protected JsonNode createResultJson(boolean useLatestTs, JsonNode result) {
long latestTs = getLatestTimestamp();
if (useLatestTs && latestTs != -1) {
ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", latestTs);
resultNode.set("values", result);
return resultNode;
} else {
return result;
}
}
}

2
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/AvgAggEntry.java

@ -35,7 +35,7 @@ public class AvgAggEntry extends BaseAggEntry {
@Override
protected double prepareResult() {
return sum.divide(BigDecimal.valueOf(count), 2, RoundingMode.HALF_UP).doubleValue();
return sum.divide(BigDecimal.valueOf(count), 10, RoundingMode.HALF_UP).doubleValue();
}
@Override

4
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/new_agg.json

@ -7,8 +7,8 @@
"allEnabledUntil": 1769907492297
},
"entityId": {
"entityType": "ASSET_PROFILE",
"id": "2b759c60-a8f4-11f0-be29-7fa922118588"
"entityType": "ASSET",
"id": "f8ad0800-a9a6-11f0-bbe6-459b63b420fe"
},
"configuration": {
"type": "LATEST_VALUES_AGGREGATION",

4
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java

@ -34,7 +34,7 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleEntityArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.LatestValuesAggregationCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
@ -59,7 +59,7 @@ public class CalculatedFieldArgumentUtils {
if (kvEntry.isPresent() && kvEntry.get().getValue() != null) {
return ArgumentEntry.createAggSingleArgument(entityId, kvEntry.get());
} else {
return new AggSingleArgumentEntry();
return new AggSingleEntityArgumentEntry();
}
}

10
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java

@ -47,7 +47,7 @@ import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.AggSingleEntityArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.LatestValuesAggregationCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmRuleState;
@ -236,7 +236,7 @@ public class CalculatedFieldUtils {
LatestValuesAggregationCalculatedFieldState aggState = (LatestValuesAggregationCalculatedFieldState) state;
Map<String, Map<EntityId, ArgumentEntry>> arguments = new HashMap<>();
proto.getAggArgumentsList().forEach(argProto -> {
AggSingleArgumentEntry entry = fromAggSingleValueArgumentProto(argProto);
AggSingleEntityArgumentEntry entry = fromAggSingleValueArgumentProto(argProto);
arguments.computeIfAbsent(argProto.getValue().getArgName(), name -> new HashMap<>()).put(entry.getEntityId(), entry);
});
arguments.forEach((argName, entityInputs) -> {
@ -248,14 +248,14 @@ public class CalculatedFieldUtils {
return state;
}
public static AggSingleArgumentEntry fromAggSingleValueArgumentProto(AggSingleArgumentEntryProto proto) {
public static AggSingleEntityArgumentEntry fromAggSingleValueArgumentProto(AggSingleArgumentEntryProto proto) {
if (!proto.hasValue()) {
return new AggSingleArgumentEntry();
return new AggSingleEntityArgumentEntry();
}
EntityId entityId = ProtoUtils.fromProto(proto.getEntityId());
SingleValueArgumentProto singleValueArgument = proto.getValue();
TsValueProto tsValueProto = singleValueArgument.getValue();
return new AggSingleArgumentEntry(
return new AggSingleEntityArgumentEntry(
entityId,
tsValueProto.getTs(),
(BasicKvEntry) KvProtoUtil.fromTsValueProto(singleValueArgument.getArgName(), tsValueProto),

344
application/src/test/java/org/thingsboard/server/cf/LatestValuesAggregationCalculatedFieldTest.java

@ -58,6 +58,7 @@ import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.cf.CalculatedFieldIntegrationTest.POLL_INTERVAL;
@DaoSqlTest
@ -116,7 +117,7 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
}
@Test
public void testNoTelemetryOnDevices_checkDefaultValueUsed() throws Exception {
public void testCreateCfOnProfile_checkInitialAggregation() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
Device device4 = createDevice("Device 4", "1234567890444");
@ -124,22 +125,153 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
createOccupancyCF("Occupied spaces", asset2.getId());
createOccupancyCF(assetProfile.getId());
await().alias("create CF and perform aggregation with default values").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "1",
"totalSpaces", "2"
));
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
postTelemetry(device3.getId(), "{\"occupied\":true}");
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "1",
"totalSpaces", "2"
));
});
}
@Test
public void testAddEntityToProfile_checkAggregation() throws Exception {
createOccupancyCF(assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
Device device4 = createDevice("Device 4", "1234567890444");
postTelemetry(device3.getId(), "{\"occupied\":true}");
postTelemetry(device4.getId(), "{\"occupied\":true}");
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
await().alias("add entity to profile with no related entities and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset2.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy).isNotNull();
assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancy.get("freeSpaces").get(0).get("value").isNull()).isTrue();
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").isNull()).isTrue();
assertThat(occupancy.get("totalSpaces").get(0).get("value").isNull()).isTrue();
});
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
await().alias("create relations and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "0",
"occupiedSpaces", "2",
"totalSpaces", "2"
));
});
postTelemetry(device3.getId(), "{\"occupied\":false}");
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "1",
"totalSpaces", "2"
));
});
}
@Test
public void testUpdateTelemetry_checkMetricsCalculation() throws Exception {
createOccupancyCF("Occupied spaces", asset.getId());
public void testChangeEntityProfile_checkAggregation() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
Device device4 = createDevice("Device 4", "1234567890444");
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
createOccupancyCF(assetProfile.getId());
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "1",
"totalSpaces", "2"
));
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
AssetProfile newAssetProfile = createAssetProfile("New Asset Profile");
asset2.setAssetProfileId(newAssetProfile.getId());
doPost("/api/asset", asset2, Asset.class);
postTelemetry(device3.getId(), "{\"occupied\":true}");
await().alias("change profile and no aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
}
@Test
public void testCreateCfOnAssetAndNoTelemetryOnDevices_checkDefaultValueUsed() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
Device device4 = createDevice("Device 4", "1234567890444");
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
createOccupancyCF(asset2.getId());
await().alias("create CF and perform aggregation with default values").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
}
@Test
public void testCreateCfAndUpdateTelemetry_checkAggregation() throws Exception {
createOccupancyCF(asset.getId());
checkInitialCalculation();
postTelemetry(device1.getId(), "{\"occupied\":false}");
@ -147,17 +279,38 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy).isNotNull();
assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
}
@Test
public void testUpdateTelemetry_checkMetricsCalculationNotExecutedUntilDeduplicationInterval() throws Exception {
createOccupancyCF("Occupied spaces", asset.getId());
public void testDeleteCf_checkNoAggregation() throws Exception {
CalculatedField cf = createOccupancyCF(asset.getId());
checkInitialCalculation();
doDelete("/api/calculatedField/" + cf.getId().getId().toString())
.andExpect(status().isOk());
postTelemetry(device1.getId(), "{\"occupied\":false}");
await().alias("delete cf and update telemetry and no aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "1",
"totalSpaces", "2"
));
});
}
@Test
public void testUpdateTelemetry_checkAggregationNotExecutedUntilDeduplicationInterval() throws Exception {
createOccupancyCF(asset.getId());
checkInitialCalculation();
postTelemetry(device1.getId(), "{\"occupied\":false}");
@ -171,136 +324,92 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
await().alias("create CF and perform initial calculation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy).isNotNull();
assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
}
@Test
public void testCreateRelation_checkMetricsCalculation() throws Exception {
createOccupancyCF("Occupied spaces", asset.getId());
checkInitialCalculation();
public void testDeleteTelemetry_checkAggregationWithPreviousValuesOrDefault() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
Device device4 = createDevice("Device 4", "1234567890444");
Device device3 = createDevice("Device 3", deviceProfile.getId(), "1234567890333");
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
postTelemetry(device3.getId(), "{\"occupied\":false}");
postTelemetry(device4.getId(), "{\"occupied\":true}");
postTelemetry(device3.getId(), "{\"occupied\":true}");
createEntityRelation(asset.getId(), device3.getId(), "Contains");
createOccupancyCF(asset2.getId());
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
await().alias("create CF and perform aggregation with default values").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy).isNotNull();
assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("3");
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "0",
"occupiedSpaces", "2",
"totalSpaces", "2"
));
});
}
@Test
public void testDeleteRelation_checkMetricsCalculation() throws Exception {
createOccupancyCF("Occupied spaces", asset.getId());
checkInitialCalculation();
deleteEntityRelation(new EntityRelation(asset.getId(), device1.getId(), "Contains", RelationTypeGroup.COMMON));
doDelete("/api/plugins/telemetry/DEVICE/" + device3.getId() + "/timeseries/delete?keys=occupied&deleteAllDataForKeys=false&rewriteLatestIfDeleted=true&deleteLatest=true&startTs=0&endTs=" + System.currentTimeMillis(), String.class);
doDelete("/api/plugins/telemetry/DEVICE/" + device4.getId() + "/timeseries/delete?keys=occupied&deleteAllDataForKeys=false&rewriteLatestIfDeleted=true&deleteLatest=true&startTs=0&endTs=" + System.currentTimeMillis(), String.class);
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
await().alias("delete latest telemetry and perform aggregation with previous or default values").atMost(deduplicationInterval * 2, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy).isNotNull();
assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("1");
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "2",
"occupiedSpaces", "0",
"totalSpaces", "2"
));
});
}
@Test
public void testCfOnProfile_checkMetricsCalculation() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
public void testCreateRelation_checkAggregation() throws Exception {
createOccupancyCF(asset.getId());
checkInitialCalculation();
Device device3 = createDevice("Device 3", deviceProfile.getId(), "1234567890333");
postTelemetry(device3.getId(), "{\"occupied\":false}");
Device device4 = createDevice("Device 4", deviceProfile.getId(), "1234567890444");
postTelemetry(device4.getId(), "{\"occupied\":false}");
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
postTelemetry(device3.getId(), "{\"occupied\":true}");
createOccupancyCF("Occupied spaces 2", assetProfile.getId());
createEntityRelation(asset.getId(), device3.getId(), "Contains");
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancyAsset1 = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancyAsset1).isNotNull();
assertThat(occupancyAsset1.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancyAsset1.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancyAsset1.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
ObjectNode occupancyAsset2 = getLatestTelemetry(asset2.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancyAsset2).isNotNull();
assertThat(occupancyAsset2.get("freeSpaces").get(0).get("value").asText()).isEqualTo("2");
assertThat(occupancyAsset2.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
assertThat(occupancyAsset2.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "2",
"totalSpaces", "3"
));
});
}
postTelemetry(device3.getId(), "{\"occupied\":true}");
@Test
public void testDeleteRelation_checkAggregation() throws Exception {
createOccupancyCF(asset.getId());
checkInitialCalculation();
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.MILLISECONDS)
deleteEntityRelation(new EntityRelation(asset.getId(), device1.getId(), "Contains", RelationTypeGroup.COMMON));
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy2 = getLatestTelemetry(asset2.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
assertThat(occupancy2).isNotNull();
assertThat(occupancy2.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancy2.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("1");
assertThat(occupancy2.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "0",
"totalSpaces", "1"
));
});
}
//
// @Test
// public void testChangeProfile_checkMetricsCalculation() throws Exception {
// DeviceProfile deviceProfile2 = doPost("/api/deviceProfile", createDeviceProfile("Device Profile 2"), DeviceProfile.class);
// device1.setDeviceProfileId(deviceProfile2.getId());
// device1 = doPost("/api/device?accessToken=" + accessToken1, device1, Device.class);
//
// postTelemetry(device1.getId(), "{\"occupied\":false}");
//
// await().alias("change profile and perform aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
// .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
// .untilAsserted(() -> {
// ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
// assertThat(occupancy).isNotNull();
// assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
// assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("0");
// assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("1");
// });
// }
//
// @Test
// public void testCfWithoutTargetProfileSpecified_checkMetricsCalculation() throws Exception {
// Device device3 = createDevice("Device 3", "1234567890333");
// postTelemetry(device3.getId(), "{\"occupied\":true}");
// createEntityRelation(asset.getId(), device3.getId(), "Contains");
//
// var configuration = (LatestValuesAggregationCalculatedFieldConfiguration) calculatedField.getConfiguration();
// configuration.getSource().setEntityProfiles(Collections.emptyList());
// calculatedField.setConfiguration(configuration);
// saveCalculatedField(calculatedField);
//
// await().alias("update cf and perform aggregation for 3 devices").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
// .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
// .untilAsserted(() -> {
// ObjectNode occupancy = getLatestTelemetry(asset.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
// assertThat(occupancy).isNotNull();
// assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1");
// assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("2");
// assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("3");
// });
// }
private void checkInitialCalculation() {
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.MILLISECONDS)
@ -316,7 +425,7 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2");
}
private CalculatedField createOccupancyCF(String name, EntityId entityId) {
private CalculatedField createOccupancyCF(EntityId entityId) {
Map<String, Argument> arguments = new HashMap<>();
Argument argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("occupied", ArgumentType.TS_LATEST, null));
@ -344,8 +453,9 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
Output output = new Output();
output.setType(OutputType.TIME_SERIES);
output.setDecimalsByDefault(0);
return createAggCf(name, entityId,
return createAggCf("Occupied spaces", entityId,
new RelationPathLevel(EntitySearchDirection.FROM, "Contains"),
arguments,
aggMetrics,
@ -393,6 +503,12 @@ public class LatestValuesAggregationCalculatedFieldTest extends AbstractControll
return doPost("/api/asset", asset, Asset.class);
}
private void verifyTelemetry(EntityId entityId, Map<String, String> expectedResults) throws Exception {
ObjectNode result = getLatestTelemetry(entityId, expectedResults.keySet().toArray(new String[0]));
assertThat(result).isNotNull();
expectedResults.forEach((key, value) -> assertThat(result.get(key).get(0).get("value").asText()).isEqualTo(value));
}
private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception {
return doGetAsync("/api/plugins/telemetry/" + entityId.getEntityType() + "/" + entityId.getId() + "/values/timeseries?keys=" + String.join(",", keys), ObjectNode.class);
}

1
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/LatestValuesAggregationCalculatedFieldConfiguration.java

@ -32,6 +32,7 @@ public class LatestValuesAggregationCalculatedFieldConfiguration implements Argu
private long deduplicationIntervalMillis;
private Map<String, AggMetric> metrics;
private Output output;
private boolean useLatestTs;
@Override
public CalculatedFieldType getType() {

Loading…
Cancel
Save