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 00e1dd3f96..a9d7c0da18 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 @@ -305,10 +305,8 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware private void onRelationChangedEvent(ComponentLifecycleMsg msg, TbCallback callback) { Function> relationAction = switch (msg.getEvent()) { - case RELATION_UPDATED -> - relatedId -> (entityId, ctx, cb) -> initRelatedEntity(entityId, relatedId, ctx, cb); - case RELATION_DELETED -> - relatedId -> (entityId, ctx, cb) -> deleteRelatedEntity(entityId, relatedId, ctx, cb); + case RELATION_UPDATED -> relatedId -> (entityId, ctx, cb) -> initRelatedEntity(entityId, relatedId, ctx, cb); + case RELATION_DELETED -> relatedId -> (entityId, ctx, cb) -> deleteRelatedEntity(entityId, relatedId, ctx, cb); default -> null; }; 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 48bce35b3e..c35d0efd5e 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 @@ -127,15 +127,6 @@ public abstract class AbstractCalculatedFieldProcessingService { return futures; } - private Map> getEntityArgumentsDuringInterval(CalculatedFieldCtx ctx, EntityId entityId, long ts) { - Map> futures = new HashMap<>(); - for (var entry : ctx.getArguments().entrySet()) { - var argValueFuture = fetchArgumentValue(ctx.getTenantId(), entityId, entry.getValue(), ts); - futures.put(entry.getKey(), argValueFuture); - } - return futures; - } - protected EntityId resolveEntityId(TenantId tenantId, EntityId entityId, Argument argument) { if (argument.getRefEntityId() != null) { return argument.getRefEntityId(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index c6a66d100e..10ace6b8cf 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -92,10 +92,8 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF @Override public Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId) { return switch (ctx.getCfType()) { - case GEOFENCING -> - resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true, System.currentTimeMillis())); - case PROPAGATION -> - resolveArgumentFutures(Map.of(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId))); + case GEOFENCING -> resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true, System.currentTimeMillis())); + case PROPAGATION -> resolveArgumentFutures(Map.of(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId))); default -> Collections.emptyMap(); }; } 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 c0208825a8..92895617c1 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 @@ -206,21 +206,6 @@ public class CalculatedFieldCtx implements Closeable { this.maxStateSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; this.maxSingleValueArgumentSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; } -// -// public boolean isRequiresScheduledReevaluation() { -// if (CalculatedFieldType.ENTITY_AGGREGATION.equals(calculatedField.getType())) { -// var configuration = (EntityAggregationCalculatedFieldConfiguration) calculatedField.getConfiguration(); -// AggInterval interval = configuration.getInterval(); -// long delayUntilIntervalEnd = interval.getDelayUntilIntervalEnd(); -// if (lastReevaluationTs < System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(systemContext.getCfCheckInterval())) { -// -// } -// if (TimeUnit.SECONDS.toMillis(systemContext.getCfCheckInterval()) >= delayUntilIntervalEnd) { -// return true; -// } -// } -// return requiresScheduledReevaluation; -// } public void init() { switch (cfType) { 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 33a8003e03..93d83476cc 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 @@ -35,8 +35,10 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; +import java.util.ArrayList; import java.util.Comparator; import java.util.HashMap; +import java.util.List; import java.util.Map; public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { @@ -117,11 +119,13 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt long now = System.currentTimeMillis(); Map> results = new HashMap<>(); + List expiredIntervals = new ArrayList<>(); for (Map.Entry> entry : intervals.entrySet()) { AggIntervalEntry intervalEntry = entry.getKey(); Map args = entry.getValue(); - processInterval(now, intervalEntry, args, results); + processInterval(now, intervalEntry, args, expiredIntervals, results); } + expiredIntervals.forEach(intervals::remove); ArrayNode result = toResult(results); if (result.isEmpty()) { @@ -157,14 +161,17 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt ctx.scheduleReevaluation(interval.getDelayUntilIntervalEnd(), actorCtx); } - private void processInterval(long now, AggIntervalEntry intervalEntry, Map args, + private void processInterval(long now, + AggIntervalEntry intervalEntry, + Map args, + List expiredIntervals, Map> results) throws Exception { long startTs = intervalEntry.getStartTs(); long endTs = intervalEntry.getEndTs(); if (now - endTs > watermarkDuration) { handleExpiredInterval(intervalEntry, args, results); - intervals.remove(intervalEntry); + expiredIntervals.add(intervalEntry); } else if (now - startTs >= intervalDuration) { handleActiveInterval(intervalEntry, args, results); } 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 780df5b795..710ce48ef6 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -102,12 +102,9 @@ public class CalculatedFieldUtils { state.getArguments().forEach((argName, argEntry) -> { switch (argEntry.getType()) { - case SINGLE_VALUE -> - builder.addSingleValueArguments(toSingleValueArgumentProto(argName, (SingleValueArgumentEntry) argEntry)); - case TS_ROLLING -> - builder.addRollingValueArguments(toRollingArgumentProto(argName, (TsRollingArgumentEntry) argEntry)); - case GEOFENCING -> - builder.addGeofencingArguments(toGeofencingArgumentProto(argName, (GeofencingArgumentEntry) argEntry)); + case SINGLE_VALUE -> builder.addSingleValueArguments(toSingleValueArgumentProto(argName, (SingleValueArgumentEntry) argEntry)); + case TS_ROLLING -> builder.addRollingValueArguments(toRollingArgumentProto(argName, (TsRollingArgumentEntry) argEntry)); + case GEOFENCING -> builder.addGeofencingArguments(toGeofencingArgumentProto(argName, (GeofencingArgumentEntry) argEntry)); case RELATED_ENTITIES -> { RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry = (RelatedEntitiesArgumentEntry) argEntry; relatedEntitiesArgumentEntry.getEntityInputs() diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java index 4c21c80086..d923eb5d7a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; +import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; @@ -25,27 +26,31 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; property = "type" ) @JsonSubTypes({ + @JsonSubTypes.Type(value = MinInterval.class, name = "MIN"), @JsonSubTypes.Type(value = HourInterval.class, name = "HOUR"), @JsonSubTypes.Type(value = DayInterval.class, name = "DAY"), @JsonSubTypes.Type(value = WeekInterval.class, name = "WEEK"), @JsonSubTypes.Type(value = WeekSunSatInterval.class, name = "WEEK_SUN_SAT"), @JsonSubTypes.Type(value = MonthInterval.class, name = "MONTH"), @JsonSubTypes.Type(value = YearInterval.class, name = "YEAR"), - @JsonSubTypes.Type(value = CustomInterval.class, name = "CUSTOM"), - - @JsonSubTypes.Type(value = MinInterval.class, name = "MIN")// todo: delete. used only for tests + @JsonSubTypes.Type(value = CustomInterval.class, name = "CUSTOM") }) @JsonIgnoreProperties(ignoreUnknown = true) public interface AggInterval { + @JsonIgnore AggIntervalType getType(); + @JsonIgnore long getIntervalDurationMillis(); + @JsonIgnore long getCurrentIntervalStartTs(); + @JsonIgnore long getCurrentIntervalEndTs(); + @JsonIgnore long getDelayUntilIntervalEnd(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalType.java index 82cf46c6ff..96bac69128 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalType.java @@ -17,14 +17,13 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i public enum AggIntervalType { + MIN, HOUR, DAY, WEEK, WEEK_SUN_SAT, MONTH, YEAR, - CUSTOM, - - MIN// todo: delete + CUSTOM } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java index cc207133fa..d8df2d9826 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java @@ -22,6 +22,7 @@ import java.time.Duration; import java.time.Instant; import java.time.LocalDate; import java.time.LocalTime; +import java.time.ZoneId; import java.time.ZonedDateTime; import java.time.temporal.ChronoUnit; import java.time.temporal.TemporalAdjusters; @@ -29,6 +30,7 @@ import java.time.temporal.TemporalAdjusters; @Data public abstract class BaseAggInterval implements AggInterval { + protected String tz; protected long offsetMillis; // delay millis since start of interval @Override @@ -54,7 +56,8 @@ public abstract class BaseAggInterval implements AggInterval { } protected long getCurrentIntervalStartTs(AggIntervalType type, int multiplier) { - ZonedDateTime now = ZonedDateTime.now(); + ZoneId zoneId = ZoneId.of(tz); + ZonedDateTime now = ZonedDateTime.now(zoneId); ZonedDateTime shiftedNow = now.minus(Duration.ofMillis(offsetMillis)); ZonedDateTime alignedStart = getAlignedBoundary(type, multiplier, false, shiftedNow); ZonedDateTime actualStart = alignedStart.plus(Duration.ofMillis(offsetMillis)); @@ -67,7 +70,8 @@ public abstract class BaseAggInterval implements AggInterval { } protected long getCurrentIntervalEndTs(AggIntervalType type, int multiplier) { - ZonedDateTime now = ZonedDateTime.now(); + ZoneId zoneId = ZoneId.of(tz); + ZonedDateTime now = ZonedDateTime.now(zoneId); ZonedDateTime shiftedNow = now.minus(Duration.ofMillis(offsetMillis)); ZonedDateTime alignedEnd = getAlignedBoundary(type, multiplier, true, shiftedNow); ZonedDateTime actualEnd = alignedEnd.plus(Duration.ofMillis(offsetMillis)); @@ -100,7 +104,7 @@ public abstract class BaseAggInterval implements AggInterval { private ZonedDateTime alignByMin(ZonedDateTime now, int multiplier, boolean next) { ZonedDateTime startOfHour = now.truncatedTo(ChronoUnit.HOURS); - long minsSinceHour = Duration.between(startOfHour, now).toHours(); + long minsSinceHour = Duration.between(startOfHour, now).toMinutes(); long aligned = (minsSinceHour / multiplier) * multiplier; if (next) aligned += multiplier; return startOfHour.plusMinutes(aligned); 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 04a9adc01b..cc56b2aef3 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 @@ -16,13 +16,24 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +@EqualsAndHashCode(callSuper = true) @Data +@NoArgsConstructor public class CustomInterval extends BaseAggInterval { private int multiplier; // number of base units (e.g. 2 hours, 5 days) private AggIntervalType internalIntervalType; + public CustomInterval(int multiplier, AggIntervalType internalIntervalType, long offsetMillis, String tz) { + this.tz = tz; + this.offsetMillis = offsetMillis; + this.multiplier = multiplier; + this.internalIntervalType = internalIntervalType; + } + @Override public AggIntervalType getType() { return AggIntervalType.CUSTOM; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java index 01cbdf2e97..e5f48d3116 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class DayInterval extends BaseAggInterval { @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java index ce84b57ae1..dfd7b7efda 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class HourInterval extends BaseAggInterval { @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MinInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MinInterval.java index f84cdb9b9e..066bc230b8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MinInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MinInterval.java @@ -16,13 +16,15 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class MinInterval extends BaseAggInterval { @Override public AggIntervalType getType() { - return AggIntervalType.HOUR; + return AggIntervalType.MIN; } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java index fe8d60f41c..7225eaca9f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class MonthInterval extends BaseAggInterval { @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java index 52054a12c3..20d03194ed 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java @@ -15,9 +15,13 @@ */ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; +import lombok.AllArgsConstructor; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@AllArgsConstructor +@NoArgsConstructor public class Watermark { private long duration; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java index 5a93076772..2ee5d5f81c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class WeekInterval extends BaseAggInterval { @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java index c70dd79a9f..f2d403c173 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class WeekSunSatInterval extends BaseAggInterval { @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java index 3c600064d1..23aedf4932 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java @@ -16,8 +16,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.NoArgsConstructor; @Data +@NoArgsConstructor public class YearInterval extends BaseAggInterval { @Override