diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 63aaf2bd68..ce29dc06e0 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/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 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 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 geofencingArgumentNames = ctx.getLinkedEntityAndCurrentOwnerGeofencingArgumentNames(); - return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), geofencingArgumentNames, scope, removedAttrKeys); + List relatedArgumentNames = ctx.getRelatedEntityArgumentNames(); + return mapToArgumentsWithDefaultValue(entityId, argNames, ctx.getArguments(), geofencingArgumentNames, relatedArgumentNames, scope, removedAttrKeys); } private Map mapToArgumentsWithDefaultValue(CalculatedFieldCtx ctx, AttributeScopeProto scope, List 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 mapToArgumentsWithDefaultValue(Map argNames, Map configArguments, List geofencingArgNames, AttributeScopeProto scope, List removedAttrKeys) { + private Map mapToArgumentsWithDefaultValue(EntityId msgEntityId, + Map argNames, + Map configArguments, + List geofencingArgNames, + List relatedEntityArgNames, + AttributeScopeProto scope, + List removedAttrKeys) { Map 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 mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, List removedTelemetryKeys) { + private Map mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, EntityId entityId, List removedTelemetryKeys) { Map deletedArguments = ctx.getArguments().entrySet().stream() .filter(entry -> removedTelemetryKeys.contains(entry.getValue().getRefEntityKey().getKey())) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); Map 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; } diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 37ec4cc3d1..af5a672157 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -516,12 +516,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware .toList(); } - private List getCfsWithRelationToEntity(EntityId entityId) { - return aggCalculatedFields.values().stream() - .filter(cf -> !findRelationsForCf(entityId, cf).isEmpty()) - .toList(); - } - private List findRelationsForCf(EntityId entityId, CalculatedFieldCtx cf) { List result = new ArrayList<>(); if (cf.getCalculatedField().getConfiguration() instanceof LatestValuesAggregationCalculatedFieldConfiguration configuration) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index c30a12a8f3..94695bcc8c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/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 fetchAggArgumentEntry(TenantId tenantId, List aggEntities, Argument argument, long startTs) {List>> futures = aggEntities.stream() + public ListenableFuture fetchAggArgumentEntry(TenantId tenantId, List aggEntities, Argument argument, long startTs) { + List>> futures = aggEntities.stream() .map(entityId -> { ListenableFuture singleAggEntryFut = fetchSingleAggArgumentEntry(tenantId, entityId, argument, startTs); return Futures.transform(singleAggEntryFut, singleAggEntry -> Map.entry(entityId, singleAggEntry), MoreExecutors.directExecutor()); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java index 4e3d00ee62..ca0a12e1d2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java +++ b/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); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index f75711a107..2d853fd3fd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/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); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index f8b9f5b665..3a9d75b3d4 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/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 mainEntityGeofencingArgumentNames; private List linkedEntityAndCurrentOwnerGeofencingArgumentNames; + private List 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; } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 65cb595632..8886482ef8 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/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) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java index 288b486e83..e81201c961 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java +++ b/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()) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggArgumentEntry.java index 12ae2c4638..7e3a8623e4 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggArgumentEntry.java +++ b/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 { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleEntityArgumentEntry.java similarity index 79% rename from application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleArgumentEntry.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleEntityArgumentEntry.java index 6b81d5380c..32ce77311d 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/AggSingleArgumentEntry.java +++ b/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; } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/LatestValuesAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/LatestValuesAggregationCalculatedFieldState.java index d25cde020c..1df27b4cbd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/LatestValuesAggregationCalculatedFieldState.java +++ b/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 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 performCalculation(Map 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 entry : metrics.entrySet()) { - String metricKey = entry.getKey(); - AggMetric metric = entry.getValue(); - - AggEntry aggMetric = AggFunctionFactory.createAggFunction(metric.getFunction()); - - for (Map 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 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 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; + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/AvgAggEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/AvgAggEntry.java index ad1f2ee8a8..afe6abb93e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/AvgAggEntry.java +++ b/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 diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/new_agg.json b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/new_agg.json index 32cd053d41..c6b841b673 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/function/new_agg.json +++ b/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", diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java index 6485508602..74cdb0cddb 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/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(); } } diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java index d4ee7bb2e3..19ad4bfb7c 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/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> 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), diff --git a/application/src/test/java/org/thingsboard/server/cf/LatestValuesAggregationCalculatedFieldTest.java b/application/src/test/java/org/thingsboard/server/cf/LatestValuesAggregationCalculatedFieldTest.java index cd9c1c578b..a4b43e362b 100644 --- a/application/src/test/java/org/thingsboard/server/cf/LatestValuesAggregationCalculatedFieldTest.java +++ b/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 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 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); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/LatestValuesAggregationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/LatestValuesAggregationCalculatedFieldConfiguration.java index 89f7e9a339..43a3360ce6 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/LatestValuesAggregationCalculatedFieldConfiguration.java +++ b/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 metrics; private Output output; + private boolean useLatestTs; @Override public CalculatedFieldType getType() {