diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntry.java index acbaffd5d9..338e667dd2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntry.java @@ -29,4 +29,8 @@ public class AggIntervalEntry { return ts >= startTs && ts < endTs; } + public long getIntervalDuration() { + return endTs - startTs; + } + } 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 54f094c289..a1208c5054 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 @@ -37,6 +37,9 @@ 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.time.Instant; +import java.time.ZoneId; +import java.time.ZonedDateTime; import java.util.ArrayList; import java.util.Comparator; import java.util.HashMap; @@ -49,7 +52,6 @@ import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDe public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { private AggInterval interval; - private long intervalDuration; private long watermarkDuration; private long checkInterval; private Map metrics; @@ -65,7 +67,6 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt super.setCtx(ctx, actorCtx); this.cfProcessingService = ctx.getCfProcessingService(); var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); - intervalDuration = configuration.getInterval().getIntervalDurationMillis(); Watermark watermark = configuration.getWatermark(); watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); checkInterval = TimeUnit.SECONDS.toMillis(ctx.getCfCheckInterval()); @@ -127,18 +128,21 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt } public void fillMissingIntervals() { + ZoneId zoneId = interval.getZoneId(); long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); - long intervalDuration = interval.getIntervalDurationMillis(); + Map> intervals = getIntervals(); AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); if (lastIntervalEntry == null) { return; } - long nextStartTs = lastIntervalEntry.getEndTs(); - long nextEndTs = nextStartTs + intervalDuration; + ZonedDateTime nextStart = Instant.ofEpochMilli(lastIntervalEntry.getEndTs()).atZone(zoneId); + ZonedDateTime nextEnd = interval.getNextIntervalStart(nextStart); - while (nextEndTs <= currentIntervalEndTs) { + while (nextEnd.toInstant().toEpochMilli() <= currentIntervalEndTs) { + long nextStartTs = nextStart.toInstant().toEpochMilli(); + long nextEndTs = nextEnd.toInstant().toEpochMilli(); AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs); arguments.forEach((argName, argumentEntry) -> { @@ -147,8 +151,8 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt entityAggEntry.getAggIntervals().computeIfAbsent(missing, missingInterval -> intervalEntryStatus); }); - nextStartTs = nextEndTs; - nextEndTs += intervalDuration; + nextStart = nextEnd; + nextEnd = interval.getNextIntervalStart(nextStart); } } @@ -174,7 +178,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt if (now - endTs > watermarkDuration) { handleExpiredInterval(intervalEntry, args, results); expiredIntervals.add(intervalEntry); - } else if (now - startTs >= intervalDuration) { + } else if (now - startTs >= intervalEntry.getIntervalDuration()) { handleActiveInterval(intervalEntry, args, results); } } diff --git a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java index fb92ebe3af..7a8f5bb625 100644 --- a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -95,7 +95,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCf_checkAggregation() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval(30L, 0L, "Europe/Kyiv"); + CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); @@ -108,7 +108,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); - long interval = customInterval.getIntervalDurationMillis(); + long interval = customInterval.getCurrentIntervalDurationMillis(); Watermark watermark = new Watermark(60); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); @@ -126,7 +126,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCf_checkAggregationDuringWatermark() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval(30L, 0L, "Europe/Kyiv"); + CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); @@ -139,7 +139,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); - long interval = customInterval.getIntervalDurationMillis(); + long interval = customInterval.getCurrentIntervalDurationMillis(); Watermark watermark = new Watermark(60); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); 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 b534002885..6bff7b4398 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 @@ -20,6 +20,9 @@ import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; +import java.time.ZoneId; +import java.time.ZonedDateTime; + @JsonTypeInfo( use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.PROPERTY, @@ -42,15 +45,21 @@ public interface AggInterval { AggIntervalType getType(); @JsonIgnore - long getIntervalDurationMillis(); + ZoneId getZoneId(); + + @JsonIgnore + long getCurrentIntervalDurationMillis(); @JsonIgnore long getCurrentIntervalStartTs(); + long getDateTimeIntervalStartTs(ZonedDateTime dateTime); + @JsonIgnore long getCurrentIntervalEndTs(); - @JsonIgnore - long getDelayUntilIntervalEnd(); + long getDateTimeIntervalEndTs(ZonedDateTime dateTime); + + ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart); } 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 f7cae9cd9e..53f7117bfc 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 @@ -15,54 +15,50 @@ */ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; -import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonInclude; -import jakarta.validation.constraints.Min; import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; +import lombok.AllArgsConstructor; import lombok.Data; +import lombok.NoArgsConstructor; -import java.time.DayOfWeek; -import java.time.Duration; -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; @Data @JsonInclude(JsonInclude.Include.NON_NULL) +@AllArgsConstructor +@NoArgsConstructor public abstract class BaseAggInterval implements AggInterval { @NotBlank protected String tz; protected Long offsetSec; // delay seconds since start of interval - @JsonIgnore - protected long getOffsetSec() { + @Override + public ZoneId getZoneId() { + return ZoneId.of(tz); + } + + protected long getOffset() { return offsetSec != null ? offsetSec : 0L; } @Override - public long getIntervalDurationMillis() { - return switch (getType()) { - case HOUR -> Duration.ofHours(1).toMillis(); - case DAY -> Duration.ofDays(1).toMillis(); - case WEEK, WEEK_SUN_SAT -> Duration.ofDays(7L).toMillis(); - case MONTH -> Duration.ofDays(Math.round(30)).toMillis(); // average - case QUARTER -> Duration.ofDays(Math.round(91)).toMillis(); - case YEAR -> Duration.ofDays(Math.round(365)).toMillis(); - default -> throw new IllegalArgumentException("Unsupported type: " + getType()); - }; + public long getCurrentIntervalDurationMillis() { + return getCurrentIntervalEndTs() - getCurrentIntervalStartTs(); } @Override public long getCurrentIntervalStartTs() { - ZoneId zoneId = ZoneId.of(tz); + ZoneId zoneId = getZoneId(); ZonedDateTime now = ZonedDateTime.now(zoneId); - long offset = getOffsetSec(); - ZonedDateTime shiftedNow = now.minusSeconds(offset); + return getDateTimeIntervalStartTs(now); + } + + @Override + public long getDateTimeIntervalStartTs(ZonedDateTime dateTime) { + long offset = getOffset(); + ZonedDateTime shiftedNow = dateTime.minusSeconds(offset); ZonedDateTime alignedStart = getAlignedBoundary(shiftedNow, false); ZonedDateTime actualStart = alignedStart.plusSeconds(offset); return actualStart.toInstant().toEpochMilli(); @@ -70,72 +66,20 @@ public abstract class BaseAggInterval implements AggInterval { @Override public long getCurrentIntervalEndTs() { - ZoneId zoneId = ZoneId.of(tz); + ZoneId zoneId = getZoneId(); ZonedDateTime now = ZonedDateTime.now(zoneId); - long offset = getOffsetSec(); - ZonedDateTime shiftedNow = now.minusSeconds(offset); - ZonedDateTime alignedEnd = getAlignedBoundary(shiftedNow, true); - ZonedDateTime actualEnd = alignedEnd.plusSeconds(offset); - return actualEnd.toInstant().toEpochMilli(); + return getDateTimeIntervalEndTs(now); } @Override - public long getDelayUntilIntervalEnd() { - long currentIntervalEndTs = getCurrentIntervalEndTs(); - long now = System.currentTimeMillis(); - return currentIntervalEndTs - now; - } - - protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { - return switch (getType()) { - case HOUR -> alignByHours(reference, next); - case DAY -> alignByDays(reference, next); - case WEEK -> alignByWeeks(reference, DayOfWeek.MONDAY, next); - case WEEK_SUN_SAT -> alignByWeeks(reference, DayOfWeek.SUNDAY, next); - case MONTH -> alignByMonths(reference, next); - case QUARTER -> alignByQuarters(reference, next); - case YEAR -> alignByYears(reference, next); - default -> throw new IllegalArgumentException("Unsupported interval type: " + getType()); - }; - } - - private ZonedDateTime alignByHours(ZonedDateTime now, boolean next) { - ZonedDateTime base = now.truncatedTo(ChronoUnit.HOURS); - return next ? base.plusHours(1) : base; - } - - private ZonedDateTime alignByDays(ZonedDateTime now, boolean next) { - ZonedDateTime base = now.truncatedTo(ChronoUnit.DAYS); - return next ? base.plusDays(1) : base; - } - - private ZonedDateTime alignByWeeks(ZonedDateTime now, DayOfWeek startOfWeek, boolean next) { - ZonedDateTime startOfWeekDate = now.with(TemporalAdjusters.previousOrSame(startOfWeek)) - .truncatedTo(ChronoUnit.DAYS); - return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate; - } - - private ZonedDateTime alignByMonths(ZonedDateTime now, boolean next) { - ZonedDateTime base = now.withDayOfMonth(1).truncatedTo(ChronoUnit.DAYS); - return next ? base.plusMonths(1) : base; - } - - private ZonedDateTime alignByQuarters(ZonedDateTime now, boolean next) { - int month = now.getMonthValue(); - int quarterStartMonth = ((month - 1) / 3) * 3 + 1; // 1, 4, 7, 10 - ZonedDateTime base = ZonedDateTime.of( - LocalDate.of(now.getYear(), quarterStartMonth, 1), - LocalTime.MIDNIGHT, - now.getZone()); - return next ? base.plusMonths(3) : base; + public long getDateTimeIntervalEndTs(ZonedDateTime dateTime) { + long offset = getOffset(); + ZonedDateTime shiftedNow = dateTime.minusSeconds(offset); + ZonedDateTime alignedEnd = getAlignedBoundary(shiftedNow, true); + ZonedDateTime actualEnd = alignedEnd.plusSeconds(offset); + return actualEnd.toInstant().toEpochMilli(); } - private ZonedDateTime alignByYears(ZonedDateTime now, boolean next) { - ZonedDateTime base = ZonedDateTime.of( - LocalDate.of(now.getYear(), 1, 1), - LocalTime.MIDNIGHT, - now.getZone()); - return next ? base.plusYears(1) : base; - } + protected abstract ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next); } 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 c8e3ee15a6..39b2f59b15 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 @@ -20,9 +20,8 @@ import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; import java.time.Duration; -import java.time.ZoneId; +import java.time.Instant; import java.time.ZonedDateTime; -import java.util.concurrent.TimeUnit; @EqualsAndHashCode(callSuper = true) @Data @@ -31,9 +30,8 @@ public class CustomInterval extends BaseAggInterval { private Long durationSec; - public CustomInterval(Long durationSec, Long offsetMillis, String tz) { - this.tz = tz; - this.offsetSec = offsetMillis; + public CustomInterval(String tz, Long offsetSec, Long durationSec) { + super(tz, offsetSec); this.durationSec = durationSec; } @@ -43,32 +41,26 @@ public class CustomInterval extends BaseAggInterval { } @Override - public long getIntervalDurationMillis() { - return Duration.ofSeconds(durationSec).toMillis(); + public long getCurrentIntervalDurationMillis() { + return getDurationMillis(); } - @Override - public long getCurrentIntervalStartTs() { - ZoneId zoneId = ZoneId.of(tz); - ZonedDateTime now = ZonedDateTime.now(zoneId); - ZonedDateTime shiftedNow = now.minusSeconds(getOffsetSec()); - - long durationMillis = getIntervalDurationMillis(); - long shiftedNowMillis = shiftedNow.toInstant().toEpochMilli(); - long alignedStartMillis = (shiftedNowMillis / durationMillis) * durationMillis; - - long offsetMillis = TimeUnit.SECONDS.toMillis(getOffsetSec()); - return alignedStartMillis + offsetMillis; + private long getDurationMillis() { + return Duration.ofSeconds(durationSec).toMillis(); } @Override - public long getCurrentIntervalEndTs() { - return getCurrentIntervalStartTs() + getIntervalDurationMillis(); + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + long durationMillis = getDurationMillis(); + long nowMillis = reference.toInstant().toEpochMilli(); + long alignedStartMillis = (nowMillis / durationMillis) * durationMillis; + ZonedDateTime aligned = Instant.ofEpochMilli(alignedStartMillis).atZone(getZoneId()); + return next ? aligned.plusSeconds(durationSec) : aligned; } @Override - public long getDelayUntilIntervalEnd() { - return getCurrentIntervalEndTs() - System.currentTimeMillis(); + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusSeconds(durationSec); } } 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 e5f48d3116..37e75c9ee6 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 @@ -18,6 +18,9 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; + @Data @NoArgsConstructor public class DayInterval extends BaseAggInterval { @@ -27,4 +30,19 @@ public class DayInterval extends BaseAggInterval { return AggIntervalType.DAY; } + public DayInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime base = reference.truncatedTo(ChronoUnit.DAYS); + return next ? base.plusDays(1) : base; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusDays(1); + } + } 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 dfd7b7efda..1cac0017e7 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,15 +16,35 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; import lombok.Data; +import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; + +@EqualsAndHashCode(callSuper = true) @Data @NoArgsConstructor public class HourInterval extends BaseAggInterval { + public HourInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + @Override public AggIntervalType getType() { return AggIntervalType.HOUR; } + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime base = reference.truncatedTo(ChronoUnit.HOURS); + return next ? base.plusHours(1) : base; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusHours(1); + } + } 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 91fc3d0413..0a540e49cd 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 @@ -18,6 +18,9 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; + @Data @NoArgsConstructor public class MonthInterval extends BaseAggInterval { @@ -27,4 +30,19 @@ public class MonthInterval extends BaseAggInterval { return AggIntervalType.MONTH; } + public MonthInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime base = reference.withDayOfMonth(1).truncatedTo(ChronoUnit.DAYS); + return next ? base.plusMonths(1) : base; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusMonths(1); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java index eb774b2341..bd27c681f8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java @@ -18,6 +18,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.LocalDate; +import java.time.LocalTime; +import java.time.ZonedDateTime; + @Data @NoArgsConstructor public class QuarterInterval extends BaseAggInterval { @@ -27,4 +31,24 @@ public class QuarterInterval extends BaseAggInterval { return AggIntervalType.QUARTER; } + public QuarterInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + int month = reference.getMonthValue(); + int quarterStartMonth = ((month - 1) / 3) * 3 + 1; // 1, 4, 7, 10 + ZonedDateTime base = ZonedDateTime.of( + LocalDate.of(reference.getYear(), quarterStartMonth, 1), + LocalTime.MIDNIGHT, + reference.getZone()); + return next ? base.plusMonths(3) : base; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusMonths(3); + } + } 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 2ee5d5f81c..381fb3bb66 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 @@ -18,6 +18,11 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.DayOfWeek; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; +import java.time.temporal.TemporalAdjusters; + @Data @NoArgsConstructor public class WeekInterval extends BaseAggInterval { @@ -27,4 +32,20 @@ public class WeekInterval extends BaseAggInterval { return AggIntervalType.WEEK; } + public WeekInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime startOfWeekDate = reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.MONDAY)) + .truncatedTo(ChronoUnit.DAYS); + return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusWeeks(1); + } + } 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 f2d403c173..242f3b7914 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 @@ -18,6 +18,11 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.DayOfWeek; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; +import java.time.temporal.TemporalAdjusters; + @Data @NoArgsConstructor public class WeekSunSatInterval extends BaseAggInterval { @@ -27,4 +32,20 @@ public class WeekSunSatInterval extends BaseAggInterval { return AggIntervalType.WEEK_SUN_SAT; } + public WeekSunSatInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime startOfWeekDate = reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.SUNDAY)) + .truncatedTo(ChronoUnit.DAYS); + return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusWeeks(1); + } + } 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 23aedf4932..83c8f58301 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 @@ -18,6 +18,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i import lombok.Data; import lombok.NoArgsConstructor; +import java.time.LocalDate; +import java.time.LocalTime; +import java.time.ZonedDateTime; + @Data @NoArgsConstructor public class YearInterval extends BaseAggInterval { @@ -27,4 +31,22 @@ public class YearInterval extends BaseAggInterval { return AggIntervalType.YEAR; } + public YearInterval(String tz, Long offsetSec) { + super(tz, offsetSec); + } + + @Override + protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { + ZonedDateTime base = ZonedDateTime.of( + LocalDate.of(reference.getYear(), 1, 1), + LocalTime.MIDNIGHT, + reference.getZone()); + return next ? base.plusYears(1) : base; + } + + @Override + public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { + return currentStart.plusYears(1); + } + } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalTest.java new file mode 100644 index 0000000000..38fda87879 --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalTest.java @@ -0,0 +1,139 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; + +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.time.Duration; +import java.time.Instant; +import java.time.ZoneId; +import java.time.ZonedDateTime; +import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.function.LongFunction; +import java.util.stream.Stream; + +import static org.assertj.core.api.Assertions.assertThat; + +public class AggIntervalTest { + + private static final String TZ = "Europe/Kiev"; + + @ParameterizedTest + @MethodSource("intervals") + void testGetStartAndEndWithoutOffset(LongFunction intervalCreator) { + AggInterval interval = intervalCreator.apply(0L); + + ZonedDateTime dateTime = ZonedDateTime.of( + // 2025.11.11 00:00:00 + 2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ) + ); + long startTs = interval.getDateTimeIntervalStartTs(dateTime); + long endTs = interval.getDateTimeIntervalEndTs(dateTime); + + assertThat(endTs).isGreaterThan(startTs); + assertThat(endTs - startTs).isEqualTo(interval.getCurrentIntervalDurationMillis()); + } + + @ParameterizedTest + @MethodSource("intervals") + void testApplyOffset(LongFunction intervalCreator) { + long offsetSec = TimeUnit.MINUTES.toSeconds(15); + AggInterval intervalWithOffset = intervalCreator.apply(offsetSec); + AggInterval intervalNoOffset = intervalCreator.apply(0L); + + ZonedDateTime dateTime = ZonedDateTime.of( + // 2025.11.11 11:20:00 - chosen so 15m offset shifts into a new interval + 2025, 11, 11, 11, 20, 0, 0, ZoneId.of(TZ) + ); + + long startWithOffsetTs = intervalWithOffset.getDateTimeIntervalStartTs(dateTime); + long startNoOffsetTs = intervalNoOffset.getDateTimeIntervalStartTs(dateTime); + + ZonedDateTime startWithOffset = Instant.ofEpochMilli(startWithOffsetTs).atZone(intervalWithOffset.getZoneId()); + ZonedDateTime startNoOffset = Instant.ofEpochMilli(startNoOffsetTs).atZone(intervalNoOffset.getZoneId()); + + long actualOffset = Duration.between(startNoOffset, startWithOffset).toSeconds(); + assertThat(actualOffset).isEqualTo(offsetSec); + } + + private static Stream intervals() { + return Stream.of( + Arguments.of((LongFunction) offset -> new HourInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new DayInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new WeekInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new WeekSunSatInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new MonthInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new QuarterInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new YearInterval(TZ, offset)), + Arguments.of((LongFunction) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4))) + ); + } + + @ParameterizedTest + @MethodSource("nextIntervalFromExactDate") + void testNextIntervalFromExactDate(LongFunction intervalCreator, Function expectedDateTimeFunction) { + AggInterval interval = intervalCreator.apply(0L); + + ZonedDateTime currentStart = ZonedDateTime.of( + 2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ) + ); + + ZonedDateTime nextStart = interval.getNextIntervalStart(currentStart); + + assertThat(nextStart).isEqualTo(expectedDateTimeFunction.apply(currentStart)); + } + + private static Stream nextIntervalFromExactDate() { + return Stream.of( + Arguments.of( + (LongFunction) offset -> new HourInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusHours(1) + ), + Arguments.of( + (LongFunction) offset -> new DayInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusDays(1) + ), + Arguments.of( + (LongFunction) offset -> new WeekInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusWeeks(1) + ), + Arguments.of( + (LongFunction) offset -> new WeekSunSatInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusWeeks(1) + ), + Arguments.of( + (LongFunction) offset -> new MonthInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusMonths(1) + ), + Arguments.of( + (LongFunction) offset -> new QuarterInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusMonths(3) + ), + Arguments.of( + (LongFunction) offset -> new YearInterval(TZ, offset), + (Function) currentInterval -> currentInterval.plusYears(1) + ), + Arguments.of( + (LongFunction) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4)), + (Function) currentInterval -> currentInterval.plusHours(4) + ) + ); + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java index fa688a0a9e..c10da4e6c6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java @@ -138,7 +138,7 @@ public class CalculatedFieldDataValidator extends DataValidator if (minAggregationIntervalInSec <= 0) { return; } - if (aggConfiguration.getInterval().getIntervalDurationMillis() < TimeUnit.SECONDS.toMillis(minAggregationIntervalInSec)) { + if (aggConfiguration.getInterval().getCurrentIntervalDurationMillis() < TimeUnit.SECONDS.toMillis(minAggregationIntervalInSec)) { throw new IllegalArgumentException("Aggregation interval duration is less than configured " + "minimum allowed aggregation interval in tenant profile: " + minAggregationIntervalInSec + " sec."); }