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 e438274ab6..e8174967a5 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 @@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesArgumentEntry; +import org.thingsboard.server.service.cf.ctx.state.aggregation.single.EntityAggregationArgumentEntry; import org.thingsboard.server.utils.CalculatedFieldUtils; import java.io.Closeable; @@ -82,6 +83,8 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, validateNewEntry(key, newEntry); if (existingEntry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { relatedEntitiesArgumentEntry.updateEntry(newEntry); + } else if (existingEntry instanceof EntityAggregationArgumentEntry entityAggArgumentEntry) { + entityAggArgumentEntry.updateEntry(newEntry); } else { arguments.put(key, newEntry); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java index d738c7d10b..7ec5098bc3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java @@ -48,19 +48,25 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { @Override public boolean updateEntry(ArgumentEntry entry) { + boolean updated = false; if (entry instanceof EntityAggregationArgumentEntry entityAggEntry) { aggIntervals.putAll(entityAggEntry.getAggIntervals()); } else if (entry instanceof SingleValueArgumentEntry singleValueArgEntry) { long entryTs = singleValueArgEntry.getTs(); long argUpdateTs = System.currentTimeMillis(); for (Map.Entry aggIntervalEntry : aggIntervals.entrySet()) { + if (singleValueArgEntry.isForceResetPrevious()) { + aggIntervalEntry.getValue().setLastArgsRefreshTs(argUpdateTs); + updated = true; + continue; + } if (aggIntervalEntry.getKey().belongsToInterval(entryTs)) { aggIntervalEntry.getValue().setLastArgsRefreshTs(argUpdateTs); return true; } } } - return false; + return updated; } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java index 04179bdd31..fbfdd9a258 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java @@ -188,6 +188,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results) { args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); processMetric(intervalEntry, argName, false, results); } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java index 39b2f59b15..c760de1143 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotNull; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; @@ -28,6 +30,8 @@ import java.time.ZonedDateTime; @NoArgsConstructor public class CustomInterval extends BaseAggInterval { + @NotNull + @Min(1) private Long durationSec; public CustomInterval(String tz, Long offsetSec, Long durationSec) {