From a35bbcc1de6d2ddfb845ca6c8295aa7496ec15df Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 24 Dec 2025 16:31:11 +0200 Subject: [PATCH 1/2] handle future telemetry and fill intervals after restart based on watermark --- ...tractCalculatedFieldProcessingService.java | 6 +- .../EntityAggregationArgumentEntry.java | 74 ++++++++++++++++--- ...EntityAggregationCalculatedFieldState.java | 18 +++++ .../utils/CalculatedFieldArgumentUtils.java | 4 +- .../EntityAggregationCalculatedFieldTest.java | 53 ++++++++++++- 5 files changed, 135 insertions(+), 20 deletions(-) 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 ec645085e6..cd88c93a71 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 @@ -233,7 +233,7 @@ public abstract class AbstractCalculatedFieldProcessingService { return config.getArguments().entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, - entry -> fetchTimeSeries(ctx.getTenantId(), entityId, entry.getValue(), config.getInterval(), ts) + entry -> fetchTimeSeries(ctx, entityId, entry.getValue(), config.getInterval(), ts) )); } @@ -341,11 +341,11 @@ public abstract class AbstractCalculatedFieldProcessingService { return resolveArgumentValue(argKey, argumentEntryFut); } - private ListenableFuture fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { + private ListenableFuture fetchTimeSeries(CalculatedFieldCtx ctx, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { long intervalStartTs = interval.getCurrentIntervalStartTs(); long intervalEndTs = interval.getCurrentIntervalEndTs(); ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), intervalStartTs, queryEndTs, 0, 1, Aggregation.NONE); - return fetchTimeSeriesInternal(tenantId, entityId, query, timeSeries -> transformAggregationArgument(timeSeries, intervalStartTs, intervalEndTs)); + return fetchTimeSeriesInternal(ctx.getTenantId(), entityId, query, timeSeries -> transformAggregationArgument(timeSeries, intervalStartTs, intervalEndTs, ctx)); } private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { 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 7ec5098bc3..ef34b6ae8e 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 @@ -19,11 +19,18 @@ import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.tbel.TbelCfArg; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.AggInterval; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; +import java.time.Instant; +import java.time.ZonedDateTime; import java.util.Map; +import java.util.concurrent.TimeUnit; @Data public class EntityAggregationArgumentEntry implements ArgumentEntry { @@ -32,10 +39,18 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { private boolean forceResetPrevious; + private AggInterval interval; + private long watermarkDuration; + public EntityAggregationArgumentEntry(Map aggIntervals) { this.aggIntervals = aggIntervals; } + public EntityAggregationArgumentEntry(Map aggIntervals, CalculatedFieldCtx ctx) { + this(aggIntervals); + setCtx(ctx); + } + @Override public ArgumentEntryType getType() { return ArgumentEntryType.ENTITY_AGGREGATION; @@ -46,29 +61,64 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { return aggIntervals; } + public void setCtx(CalculatedFieldCtx ctx) { + var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); + interval = configuration.getInterval(); + Watermark watermark = configuration.getWatermark(); + watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); + } + @Override public boolean updateEntry(ArgumentEntry entry) { - boolean updated = false; if (entry instanceof EntityAggregationArgumentEntry entityAggEntry) { aggIntervals.putAll(entityAggEntry.getAggIntervals()); + return true; } 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; - } + long now = System.currentTimeMillis(); + if (updateExistingIntervals(singleValueArgEntry, entryTs, now)) { + return true; + } + return createNewInterval(entryTs, now); + } + return false; + } + + private boolean updateExistingIntervals(SingleValueArgumentEntry entry, long entryTs, long now) { + boolean updated = false; + + for (Map.Entry aggIntervalEntry : aggIntervals.entrySet()) { + AggIntervalEntry interval = aggIntervalEntry.getKey(); + AggIntervalEntryStatus status = aggIntervalEntry.getValue(); + if (entry.isForceResetPrevious()) { + status.setLastArgsRefreshTs(now); + updated = true; + continue; + } + if (interval.belongsToInterval(entryTs)) { + status.setLastArgsRefreshTs(now); + return true; } } + return updated; } + private boolean createNewInterval(long entryTs, long now) { + ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(entryTs), interval.getZoneId()); + + long startTs = interval.getDateTimeIntervalStartTs(zdt); + long endTs = interval.getDateTimeIntervalEndTs(zdt); + + if (now - endTs > watermarkDuration) { + return false; + } + + AggIntervalEntry newInterval = new AggIntervalEntry(startTs, endTs); + aggIntervals.computeIfAbsent(newInterval, i -> new AggIntervalEntryStatus(now)); + return true; + } + @Override public boolean isEmpty() { return aggIntervals.isEmpty(); 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 0f6e01a344..59335717fb 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 @@ -82,6 +82,15 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt interval = configuration.getInterval(); metrics = configuration.getMetrics(); produceIntermediateResult = configuration.isProduceIntermediateResult(); + setCtxToArguments(); + } + + private void setCtxToArguments() { + arguments.values().forEach(argument -> { + if (argument instanceof EntityAggregationArgumentEntry entityAggArgument) { + entityAggArgument.setCtx(ctx); + } + }); } @Override @@ -154,8 +163,10 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt } private void fillMissingIntervals() { + long now = System.currentTimeMillis(); ZoneId zoneId = interval.getZoneId(); long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); + long watermarkThresholdTs = now - watermarkDuration; Map> intervals = getIntervals(); AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); @@ -169,6 +180,13 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt while (nextEnd.toInstant().toEpochMilli() <= currentIntervalEndTs) { long nextStartTs = nextStart.toInstant().toEpochMilli(); long nextEndTs = nextEnd.toInstant().toEpochMilli(); + + if (nextEndTs < watermarkThresholdTs) { + nextStart = nextEnd; + nextEnd = interval.getNextIntervalStart(nextStart); + continue; + } + AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs); arguments.forEach((argName, argumentEntry) -> { 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 7e0701cd2d..d6f08cbc69 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -75,7 +75,7 @@ public class CalculatedFieldArgumentUtils { return new SingleValueArgumentEntry(); } - public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs) { + public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs, CalculatedFieldCtx ctx) { Map aggIntervals = new HashMap<>(); AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(startIntervalTs, endIntervalTs); if (timeSeries == null || timeSeries.isEmpty()) { @@ -83,7 +83,7 @@ public class CalculatedFieldArgumentUtils { } else { aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus(System.currentTimeMillis())); } - return new EntityAggregationArgumentEntry(aggIntervals); + return new EntityAggregationArgumentEntry(aggIntervals, ctx); } private static KvEntry createDefaultKvEntry(Argument argument) { 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 253d02ebe5..66c11a06a6 100644 --- a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -55,6 +55,8 @@ import static org.thingsboard.server.cf.CalculatedFieldIntegrationTest.POLL_INTE @DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD) public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest { + private final String TZ = "Europe/Kyiv"; + private Tenant savedTenant; @Before @@ -93,7 +95,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCfAndNoTelemetryDuringInterval_checkAggregation() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); + CustomInterval customInterval = new CustomInterval(TZ, 0L, 5L); createConsumptionCF(device.getId(), customInterval, null); long interval = customInterval.getCurrentIntervalDurationMillis(); @@ -113,7 +115,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCfWithoutWatermark_checkAggregation() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); + CustomInterval customInterval = new CustomInterval(TZ, 0L, 5L); createConsumptionCF(device.getId(), customInterval, null); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); @@ -156,7 +158,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCfWithWatermark_checkAggregationDuringWatermark() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); + CustomInterval customInterval = new CustomInterval(TZ, 0L, 5L); Watermark watermark = new Watermark(10); createConsumptionCF(device.getId(), customInterval, watermark); @@ -196,6 +198,51 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest }); } + @Test + public void testSendFutureTelemetry_checkAggregation() throws Exception { + Device device = createDevice("Device", "1234567890111"); + + CustomInterval customInterval = new CustomInterval(TZ, 0L, 2L); + createConsumptionCF(device.getId(), customInterval, null); + + long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); + + long tsBeforeInterval = currentIntervalStartTs - 1000; + long tsInInterval_1 = currentIntervalStartTs + 1000; + long tsInInterval_2 = currentIntervalStartTs + 500; + long tsInInterval_3 = currentIntervalStartTs + 200; + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsBeforeInterval)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":100}}", tsInInterval_1)); + 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.getCurrentIntervalDurationMillis(); + + await().alias("create CF -> perform aggregation after interval end") + .atMost(2 * interval, TimeUnit.MILLISECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); + assertThat(result).isNotNull(); + assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("133"); + }); + + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":500}}", currentIntervalStartTs + 4500L)); + + await().alias("update telemetry that belongs to future interval -> check aggregation ") + .atMost(3 * interval, TimeUnit.MILLISECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); + assertThat(result).isNotNull(); + assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("500"); + assertThat(result.get("consumption").get(0).get("ts").asLong()).isEqualTo(currentIntervalStartTs + 4000L); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("500"); + assertThat(result.get("avgConsumption").get(0).get("ts").asLong()).isEqualTo(currentIntervalStartTs + 4000L); + }); + } + private CalculatedField createConsumptionCF(EntityId entityId, AggInterval aggInterval, Watermark watermark) { Map arguments = new HashMap<>(); Argument argument = new Argument(); From 00ab73f114d1bb042f4e4913d4d680354aec7db7 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 29 Dec 2025 11:34:06 +0200 Subject: [PATCH 2/2] passed ctx as param to updateEntry --- ...CalculatedFieldEntityMessageProcessor.java | 15 +++++---- ...tractCalculatedFieldProcessingService.java | 6 ++-- .../service/cf/ctx/state/ArgumentEntry.java | 2 +- .../ctx/state/BaseCalculatedFieldState.java | 8 ++--- .../ctx/state/SingleValueArgumentEntry.java | 2 +- .../cf/ctx/state/TsRollingArgumentEntry.java | 2 +- .../RelatedEntitiesArgumentEntry.java | 5 +-- .../EntityAggregationArgumentEntry.java | 28 ++++++---------- ...EntityAggregationCalculatedFieldState.java | 9 ----- .../alarm/AlarmCalculatedFieldState.java | 6 ++-- .../geofencing/GeofencingArgumentEntry.java | 3 +- .../propagation/PropagationArgumentEntry.java | 3 +- .../utils/CalculatedFieldArgumentUtils.java | 4 +-- .../GeofencingValueArgumentEntryTest.java | 25 +++++++++----- .../state/PropagationArgumentEntryTest.java | 33 +++++++++++-------- .../RelatedEntitiesArgumentEntryTest.java | 15 ++++++--- .../state/SingleValueArgumentEntryTest.java | 23 ++++++++----- .../ctx/state/TsRollingArgumentEntryTest.java | 17 +++++++--- 18 files changed, 114 insertions(+), 92 deletions(-) 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 d2ab4c85ac..a24c5fb732 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 @@ -572,24 +572,24 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } private Map mapToArguments(CalculatedFieldCtx ctx, List data) { - return mapToArguments(entityId, ctx.getMainEntityArguments(), Collections.emptyMap(), data); + return mapToArguments(entityId, ctx, ctx.getMainEntityArguments(), Collections.emptyMap(), data); } private Map mapToArguments(CalculatedFieldCtx ctx, EntityId entityId, List data) { - return mapToArguments(entityId, ctx.getLinkedAndDynamicArgs(entityId), ctx.getRelatedEntityArguments(), data); + return mapToArguments(entityId, ctx, ctx.getLinkedAndDynamicArgs(entityId), ctx.getRelatedEntityArguments(), data); } - private Map mapToArguments(EntityId originator, Map> args, Map> relatedEntityArgs, List data) { + private Map mapToArguments(EntityId originator, CalculatedFieldCtx ctx, Map> args, Map> relatedEntityArgs, List data) { Map arguments = new HashMap<>(); if (!relatedEntityArgs.isEmpty() || !args.isEmpty()) { for (TsKvProto item : data) { ReferencedEntityKey key = new ReferencedEntityKey(item.getKv().getKey(), ArgumentType.TS_LATEST, null); SingleValueArgumentEntry relatedArgIncoming = new SingleValueArgumentEntry(originator, item); - mapLatest(relatedArgIncoming, relatedEntityArgs.get(key), arguments); + mapLatest(ctx, relatedArgIncoming, relatedEntityArgs.get(key), arguments); SingleValueArgumentEntry incoming = new SingleValueArgumentEntry(item); - mapLatest(incoming, args.get(key), arguments); + mapLatest(ctx, incoming, args.get(key), arguments); key = new ReferencedEntityKey(item.getKv().getKey(), ArgumentType.TS_ROLLING, null); mapRolling(item, args.get(key), arguments); @@ -598,7 +598,8 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM return arguments; } - private void mapLatest(SingleValueArgumentEntry incoming, + private void mapLatest(CalculatedFieldCtx ctx, + SingleValueArgumentEntry incoming, Set argNames, Map arguments) { if (argNames != null) { @@ -606,7 +607,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (existing == null) { return incoming; } - existing.updateEntry(incoming); + existing.updateEntry(incoming, ctx); return existing; })); } 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 cd88c93a71..ec645085e6 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 @@ -233,7 +233,7 @@ public abstract class AbstractCalculatedFieldProcessingService { return config.getArguments().entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, - entry -> fetchTimeSeries(ctx, entityId, entry.getValue(), config.getInterval(), ts) + entry -> fetchTimeSeries(ctx.getTenantId(), entityId, entry.getValue(), config.getInterval(), ts) )); } @@ -341,11 +341,11 @@ public abstract class AbstractCalculatedFieldProcessingService { return resolveArgumentValue(argKey, argumentEntryFut); } - private ListenableFuture fetchTimeSeries(CalculatedFieldCtx ctx, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { + private ListenableFuture fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { long intervalStartTs = interval.getCurrentIntervalStartTs(); long intervalEndTs = interval.getCurrentIntervalEndTs(); ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), intervalStartTs, queryEndTs, 0, 1, Aggregation.NONE); - return fetchTimeSeriesInternal(ctx.getTenantId(), entityId, query, timeSeries -> transformAggregationArgument(timeSeries, intervalStartTs, intervalEndTs, ctx)); + return fetchTimeSeriesInternal(tenantId, entityId, query, timeSeries -> transformAggregationArgument(timeSeries, intervalStartTs, intervalEndTs)); } private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { 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 dc23ffa979..fc4f8ce365 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 @@ -52,7 +52,7 @@ public interface ArgumentEntry { Object getValue(); - boolean updateEntry(ArgumentEntry entry); + boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx); boolean isEmpty(); 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 09ca35cc4b..0f648600a5 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 @@ -86,13 +86,13 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, validateNewEntry(key, newEntry); if (existingEntry instanceof RelatedEntitiesArgumentEntry || existingEntry instanceof EntityAggregationArgumentEntry) { - updateEntry(existingEntry, newEntry); + updateEntry(existingEntry, newEntry, ctx); } else { arguments.put(key, newEntry); } entryUpdated = true; } else { - entryUpdated = updateEntry(existingEntry, newEntry); + entryUpdated = updateEntry(existingEntry, newEntry, ctx); } if (entryUpdated) { @@ -111,8 +111,8 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, return updatedArguments; } - protected boolean updateEntry(ArgumentEntry existingEntry, ArgumentEntry newEntry) { - return existingEntry.updateEntry(newEntry); + protected boolean updateEntry(ArgumentEntry existingEntry, ArgumentEntry newEntry, CalculatedFieldCtx ctx) { + return existingEntry.updateEntry(newEntry, ctx); } @Override 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 4d0c4d7724..028c3c429b 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 @@ -159,7 +159,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof SingleValueArgumentEntry singleValueEntry) { if (singleValueEntry.getTs() < this.ts) { return false; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index 8cdc9ddcf9..8abddb3d4a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -100,7 +100,7 @@ public class TsRollingArgumentEntry implements ArgumentEntry, HasLatestTs { } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof TsRollingArgumentEntry tsRollingEntry) { updateTsRollingEntry(tsRollingEntry); } else if (entry instanceof SingleValueArgumentEntry singleValueEntry) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java index 219cf471ed..0a5c850e4f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java @@ -23,6 +23,7 @@ import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.HasLatestTs; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; @@ -63,7 +64,7 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { entityInputs.putAll(relatedEntitiesArgumentEntry.entityInputs); return true; @@ -74,7 +75,7 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs } ArgumentEntry argumentEntry = entityInputs.get(singleValueArgumentEntry.getEntityId()); if (argumentEntry != null) { - argumentEntry.updateEntry(singleValueArgumentEntry); + argumentEntry.updateEntry(singleValueArgumentEntry, ctx); } else { entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); } 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 ef34b6ae8e..c0a0603390 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 @@ -39,18 +39,10 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { private boolean forceResetPrevious; - private AggInterval interval; - private long watermarkDuration; - public EntityAggregationArgumentEntry(Map aggIntervals) { this.aggIntervals = aggIntervals; } - public EntityAggregationArgumentEntry(Map aggIntervals, CalculatedFieldCtx ctx) { - this(aggIntervals); - setCtx(ctx); - } - @Override public ArgumentEntryType getType() { return ArgumentEntryType.ENTITY_AGGREGATION; @@ -61,15 +53,8 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { return aggIntervals; } - public void setCtx(CalculatedFieldCtx ctx) { - var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); - interval = configuration.getInterval(); - Watermark watermark = configuration.getWatermark(); - watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); - } - @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof EntityAggregationArgumentEntry entityAggEntry) { aggIntervals.putAll(entityAggEntry.getAggIntervals()); return true; @@ -79,7 +64,7 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { if (updateExistingIntervals(singleValueArgEntry, entryTs, now)) { return true; } - return createNewInterval(entryTs, now); + return createNewInterval(entryTs, now, ctx); } return false; } @@ -104,7 +89,14 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { return updated; } - private boolean createNewInterval(long entryTs, long now) { + private boolean createNewInterval(long entryTs, long now, CalculatedFieldCtx ctx) { + if (!(ctx.getCalculatedField().getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration config)) { + return false; + } + AggInterval interval = config.getInterval(); + Watermark watermark = config.getWatermark(); + long watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); + ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(entryTs), interval.getZoneId()); long startTs = interval.getDateTimeIntervalStartTs(zdt); 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 59335717fb..99bddc374c 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 @@ -82,15 +82,6 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt interval = configuration.getInterval(); metrics = configuration.getMetrics(); produceIntermediateResult = configuration.isProduceIntermediateResult(); - setCtxToArguments(); - } - - private void setCtxToArguments() { - arguments.values().forEach(argument -> { - if (argument instanceof EntityAggregationArgumentEntry entityAggArgument) { - entityAggArgument.setCtx(ctx); - } - }); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java index 4f2cbede09..63b7cc8eb1 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java @@ -224,17 +224,17 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState { } @Override - protected boolean updateEntry(ArgumentEntry existingArgumentEntry, ArgumentEntry newArgumentEntry) { + protected boolean updateEntry(ArgumentEntry existingArgumentEntry, ArgumentEntry newArgumentEntry, CalculatedFieldCtx ctx) { if (!(existingArgumentEntry instanceof SingleValueArgumentEntry existingEntry) || !(newArgumentEntry instanceof SingleValueArgumentEntry newEntry)) { - return super.updateEntry(existingArgumentEntry, newArgumentEntry); + return super.updateEntry(existingArgumentEntry, newArgumentEntry, ctx); } if (newEntry.getTs() < existingEntry.getTs()) { if (existingEntry.isDefaultValue()) { existingEntry.setTs(newEntry.getTs()); } } - return super.updateEntry(existingEntry, newEntry); + return super.updateEntry(existingEntry, newEntry, ctx); } public void processAlarmAction(Alarm alarm, ActionType action) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java index 01c7119993..a3305ea52d 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java @@ -25,6 +25,7 @@ import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.HasLatestTs; import java.util.Map; @@ -68,7 +69,7 @@ public class GeofencingArgumentEntry implements ArgumentEntry, HasLatestTs { } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (!(entry instanceof GeofencingArgumentEntry geofencingArgumentEntry)) { throw new IllegalArgumentException("Unsupported argument entry type for geofencing argument entry: " + entry.getType()); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java index 8536c0f65f..de04d0e817 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java @@ -21,6 +21,7 @@ import org.thingsboard.script.api.tbel.TbelCfPropagationArg; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import java.util.ArrayList; import java.util.Collection; @@ -59,7 +60,7 @@ public class PropagationArgumentEntry implements ArgumentEntry { } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (!(entry instanceof PropagationArgumentEntry updated)) { throw new IllegalArgumentException("Unsupported argument entry type for propagation argument entry: " + entry.getType()); } 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 d6f08cbc69..7e0701cd2d 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -75,7 +75,7 @@ public class CalculatedFieldArgumentUtils { return new SingleValueArgumentEntry(); } - public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs, CalculatedFieldCtx ctx) { + public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs) { Map aggIntervals = new HashMap<>(); AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(startIntervalTs, endIntervalTs); if (timeSeries == null || timeSeries.isEmpty()) { @@ -83,7 +83,7 @@ public class CalculatedFieldArgumentUtils { } else { aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus(System.currentTimeMillis())); } - return new EntityAggregationArgumentEntry(aggIntervals, ctx); + return new EntityAggregationArgumentEntry(aggIntervals); } private static KvEntry createDefaultKvEntry(Argument argument) { diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingValueArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingValueArgumentEntryTest.java index 6da4bdc882..d274da2434 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingValueArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingValueArgumentEntryTest.java @@ -18,6 +18,9 @@ package org.thingsboard.server.service.cf.ctx.state; import io.hypersistence.utils.hibernate.type.json.internal.JacksonUtil; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.common.util.geo.PerimeterDefinition; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.EntityId; @@ -33,6 +36,7 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ExtendWith(MockitoExtension.class) public class GeofencingValueArgumentEntryTest { private final AssetId ZONE_1_ID = new AssetId(UUID.fromString("c0e3031c-7df1-45e4-9590-cfd621a4d714")); @@ -46,6 +50,9 @@ public class GeofencingValueArgumentEntryTest { private GeofencingArgumentEntry entry; + @Mock + private CalculatedFieldCtx ctx; + @BeforeEach void setUp() { entry = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, allowedZoneAttributeKvEntry, ZONE_2_ID, restrictedZoneAttributeKvEntry)); @@ -58,14 +65,14 @@ public class GeofencingValueArgumentEntryTest { @Test void testUpdateEntryWhenSingleEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new SingleValueArgumentEntry())) + assertThatThrownBy(() -> entry.updateEntry(new SingleValueArgumentEntry(), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for geofencing argument entry: SINGLE_VALUE"); } @Test void testUpdateEntryWhenRollingEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) + assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for geofencing argument entry: TS_ROLLING"); } @@ -74,7 +81,7 @@ public class GeofencingValueArgumentEntryTest { void testUpdateEntryWithTheSameTs() { BaseAttributeKvEntry differentValueSameTs = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 363L, 156L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueSameTs, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isFalse(); + assertThat(entry.updateEntry(updated, ctx)).isFalse(); } @Test @@ -83,7 +90,7 @@ public class GeofencingValueArgumentEntryTest { BaseAttributeKvEntry differentValueNewVersionIsNull = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, null); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsNull, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isTrue(); + assertThat(entry.updateEntry(updated, ctx)).isTrue(); assertThat(entry.getValue()).isInstanceOf(Map.class); Map value = (Map) entry.getValue(); @@ -105,7 +112,7 @@ public class GeofencingValueArgumentEntryTest { BaseAttributeKvEntry differentValueNewVersionIsSet = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 156L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsSet, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isTrue(); + assertThat(entry.updateEntry(updated, ctx)).isTrue(); assertThat(entry.getValue()).isInstanceOf(Map.class); Map value = (Map) entry.getValue(); @@ -126,7 +133,7 @@ public class GeofencingValueArgumentEntryTest { BaseAttributeKvEntry differentValueNewVersionIsSet = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 154L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsSet, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isFalse(); + assertThat(entry.updateEntry(updated, ctx)).isFalse(); } @Test @@ -134,7 +141,7 @@ public class GeofencingValueArgumentEntryTest { BaseAttributeKvEntry newTsAndTheSameValue = new BaseAttributeKvEntry(allowedZoneDataEntry, 364L, 156L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, newTsAndTheSameValue, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isTrue(); + assertThat(entry.updateEntry(updated, ctx)).isTrue(); } @Test @@ -142,7 +149,7 @@ public class GeofencingValueArgumentEntryTest { BaseAttributeKvEntry oldTsAndTheSameValue = new BaseAttributeKvEntry(allowedZoneDataEntry, 362L, 156L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, oldTsAndTheSameValue, ZONE_2_ID, restrictedZoneAttributeKvEntry)); - assertThat(entry.updateEntry(updated)).isFalse(); + assertThat(entry.updateEntry(updated, ctx)).isFalse(); } @Test @@ -150,7 +157,7 @@ public class GeofencingValueArgumentEntryTest { final AssetId NEW_ZONE_ID = new AssetId(UUID.fromString("a3eacf1a-6af3-4e9f-87c4-502bb25c7dc3")); BaseAttributeKvEntry newZone = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 156L); var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, allowedZoneAttributeKvEntry, ZONE_2_ID, restrictedZoneAttributeKvEntry, NEW_ZONE_ID, newZone)); - assertThat(entry.updateEntry(updated)).isTrue(); + assertThat(entry.updateEntry(updated, ctx)).isTrue(); } @Test diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java index 12f3e4298d..f4098dc2df 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java @@ -17,6 +17,9 @@ package org.thingsboard.server.service.cf.ctx.state; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfPropagationArg; import org.thingsboard.server.common.data.id.AssetId; @@ -31,6 +34,7 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ExtendWith(MockitoExtension.class) public class PropagationArgumentEntryTest { private final AssetId ENTITY_1_ID = new AssetId(UUID.fromString("b0a8637d-6d67-43d5-a483-c0e391afe805")); @@ -39,6 +43,9 @@ public class PropagationArgumentEntryTest { private PropagationArgumentEntry entry; + @Mock + private CalculatedFieldCtx ctx; + @BeforeEach void setUp() { List propagationEntityIds = new ArrayList<>(); @@ -68,14 +75,14 @@ public class PropagationArgumentEntryTest { @Test void testUpdateEntryWhenSingleEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new SingleValueArgumentEntry())) + assertThatThrownBy(() -> entry.updateEntry(new SingleValueArgumentEntry(), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for propagation argument entry: SINGLE_VALUE"); } @Test void testUpdateEntryWhenRollingEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) + assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for propagation argument entry: TS_ROLLING"); } @@ -85,7 +92,7 @@ public class PropagationArgumentEntryTest { var newIds = new ArrayList(List.of(ENTITY_3_ID, ENTITY_1_ID)); var updated = new PropagationArgumentEntry(newIds); - boolean changed = entry.updateEntry(updated); + boolean changed = entry.updateEntry(updated, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyElementsOf(newIds); @@ -95,7 +102,7 @@ public class PropagationArgumentEntryTest { void testUpdateEntryClearsWhenNewEntryIsEmpty() { var updatedEmpty = new PropagationArgumentEntry(List.of()); - boolean changed = entry.updateEntry(updatedEmpty); + boolean changed = entry.updateEntry(updatedEmpty, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).isEmpty(); @@ -106,7 +113,7 @@ public class PropagationArgumentEntryTest { var added = new PropagationArgumentEntry(); added.setAdded(List.of(ENTITY_3_ID)); - boolean changed = entry.updateEntry(added); + boolean changed = entry.updateEntry(added, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); @@ -118,7 +125,7 @@ public class PropagationArgumentEntryTest { var added = new PropagationArgumentEntry(); added.setAdded(List.of(ENTITY_2_ID)); - boolean changed = entry.updateEntry(added); + boolean changed = entry.updateEntry(added, ctx); assertThat(changed).isFalse(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); @@ -130,7 +137,7 @@ public class PropagationArgumentEntryTest { var removed = new PropagationArgumentEntry(); removed.setRemoved(ENTITY_2_ID); - boolean changed = entry.updateEntry(removed); + boolean changed = entry.updateEntry(removed, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID); @@ -142,7 +149,7 @@ public class PropagationArgumentEntryTest { var removed = new PropagationArgumentEntry(); removed.setRemoved(ENTITY_3_ID); - boolean changed = entry.updateEntry(removed); + boolean changed = entry.updateEntry(removed, ctx); assertThat(changed).isFalse(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); @@ -154,7 +161,7 @@ public class PropagationArgumentEntryTest { var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID)); restore.setIgnoreRemovedEntities(true); - boolean changed = entry.updateEntry(restore); + boolean changed = entry.updateEntry(restore, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); @@ -168,7 +175,7 @@ public class PropagationArgumentEntryTest { var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); restore.setIgnoreRemovedEntities(true); - boolean changed = entry.updateEntry(restore); + boolean changed = entry.updateEntry(restore, ctx); assertThat(changed).isFalse(); // expected no change, since we consider the removal of stale ids as no-op assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID); @@ -182,7 +189,7 @@ public class PropagationArgumentEntryTest { var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_3_ID)); restore.setIgnoreRemovedEntities(true); - boolean changed = entry.updateEntry(restore); + boolean changed = entry.updateEntry(restore, ctx); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_3_ID); @@ -197,7 +204,7 @@ public class PropagationArgumentEntryTest { var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID)); restore.setIgnoreRemovedEntities(true); - boolean changed = entry.updateEntry(restore); + boolean changed = entry.updateEntry(restore, ctx); assertThat(changed).isFalse(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); @@ -211,7 +218,7 @@ public class PropagationArgumentEntryTest { var restore = new PropagationArgumentEntry(List.of()); restore.setIgnoreRemovedEntities(true); - boolean changed = entry.updateEntry(restore); + boolean changed = entry.updateEntry(restore, ctx); assertThat(changed).isFalse(); // expected no change, since we consider the removal of stale ids as no-op assertThat(entry.getEntityIds()).isEmpty(); diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesArgumentEntryTest.java index cc60b249ac..c357e2c30b 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesArgumentEntryTest.java @@ -17,6 +17,9 @@ package org.thingsboard.server.service.cf.ctx.state; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; @@ -30,10 +33,14 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ExtendWith(MockitoExtension.class) public class RelatedEntitiesArgumentEntryTest { private RelatedEntitiesArgumentEntry entry; + @Mock + private CalculatedFieldCtx ctx; + private final DeviceId device1 = new DeviceId(UUID.fromString("1984e5f4-9ff0-4187-84ae-e4438bba4c8a")); private final DeviceId device2 = new DeviceId(UUID.fromString("937fc062-1a9d-438f-aa22-55a93fc908b7")); @@ -50,7 +57,7 @@ public class RelatedEntitiesArgumentEntryTest { @Test void testUpdateEntryWhenNotAggEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) + assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for aggregation argument entry: " + ArgumentEntryType.TS_ROLLING); } @@ -65,7 +72,7 @@ public class RelatedEntitiesArgumentEntryTest { device4, new SingleValueArgumentEntry(device4, new BasicTsKvEntry(ts - 60, new LongDataEntry("key", 23L), 7L)) ), false); - assertThat(entry.updateEntry(relatedEntitiesArgumentEntry)).isTrue(); + assertThat(entry.updateEntry(relatedEntitiesArgumentEntry, ctx)).isTrue(); Map aggInputs = entry.getEntityInputs(); assertThat(aggInputs.size()).isEqualTo(4); @@ -79,7 +86,7 @@ public class RelatedEntitiesArgumentEntryTest { SingleValueArgumentEntry singleEntityArgumentEntry = new SingleValueArgumentEntry(device3, new BasicTsKvEntry(ts - 50, new LongDataEntry("key", 18L), 10L)); - assertThat(entry.updateEntry(singleEntityArgumentEntry)).isTrue(); + assertThat(entry.updateEntry(singleEntityArgumentEntry, ctx)).isTrue(); Map aggInputs = entry.getEntityInputs(); assertThat(aggInputs.size()).isEqualTo(3); @@ -90,7 +97,7 @@ public class RelatedEntitiesArgumentEntryTest { void testUpdateEntryWhenSingleValueArgumentEntryPassedAndEntryByIdExist() { SingleValueArgumentEntry singleEntityArgumentEntry = new SingleValueArgumentEntry(device2, new BasicTsKvEntry(ts - 50, new LongDataEntry("key", 18L), 10L)); - assertThat(entry.updateEntry(singleEntityArgumentEntry)).isTrue(); + assertThat(entry.updateEntry(singleEntityArgumentEntry, ctx)).isTrue(); Map aggInputs = entry.getEntityInputs(); assertThat(aggInputs.size()).isEqualTo(2); diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java index 4ada355054..e2d287c778 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java @@ -17,6 +17,9 @@ package org.thingsboard.server.service.cf.ctx.state; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; import org.thingsboard.server.common.data.kv.JsonDataEntry; @@ -29,10 +32,14 @@ import java.util.Map; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ExtendWith(MockitoExtension.class) public class SingleValueArgumentEntryTest { private SingleValueArgumentEntry entry; + @Mock + private CalculatedFieldCtx ctx; + private final long ts = System.currentTimeMillis(); @BeforeEach @@ -47,48 +54,48 @@ public class SingleValueArgumentEntryTest { @Test void testUpdateEntryWhenRollingEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) + assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L), ctx)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for single value argument entry: " + ArgumentEntryType.TS_ROLLING); } @Test void testUpdateEntryWithTheSameTs() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 363L))).isFalse(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 363L), ctx)).isFalse(); } @Test void testUpdateEntryWithTheSameTsAndDifferentVersion() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 364L))).isTrue(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 364L), ctx)).isTrue(); } @Test void testUpdateEntryWhenNewVersionIsNull() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 16, new LongDataEntry("key", 13L), null))).isTrue(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 16, new LongDataEntry("key", 13L), null), ctx)).isTrue(); assertThat(entry.getValue()).isEqualTo(13L); assertThat(entry.getVersion()).isNull(); } @Test void testUpdateEntryWhenNewVersionIsGreaterThanCurrent() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 18L), 369L))).isTrue(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 18L), 369L), ctx)).isTrue(); assertThat(entry.getValue()).isEqualTo(18L); assertThat(entry.getVersion()).isEqualTo(369L); } @Test void testUpdateEntryWhenNewVersionIsLessThanCurrent() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 18L), 234L))).isFalse(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 18L), 234L), ctx)).isFalse(); } @Test void testUpdateEntryWhenValueWasNotChanged() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 11L), 364L))).isTrue(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 11L), 364L), ctx)).isTrue(); } @Test void testUpdateEntryWithOldTs() { - assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts - 10, new LongDataEntry("key", 14L), 365L))).isFalse(); + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts - 10, new LongDataEntry("key", 14L), 365L), ctx)).isFalse(); } @Test diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java index b1f8063857..94b2a7b389 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java @@ -17,6 +17,9 @@ package org.thingsboard.server.service.cf.ctx.state; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; @@ -26,10 +29,14 @@ import java.util.TreeMap; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ExtendWith(MockitoExtension.class) public class TsRollingArgumentEntryTest { private TsRollingArgumentEntry entry; + @Mock + private CalculatedFieldCtx ctx; + private final long ts = System.currentTimeMillis(); @BeforeEach @@ -51,7 +58,7 @@ public class TsRollingArgumentEntryTest { void testUpdateEntryWhenSingleValueEntryPassed() { SingleValueArgumentEntry newEntry = new SingleValueArgumentEntry(ts - 10, new DoubleDataEntry("key", 23.0), 123L); - assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.updateEntry(newEntry, ctx)).isTrue(); assertThat(entry.getTsRecords()).hasSize(4); assertThat(entry.getTsRecords().get(ts - 10)).isEqualTo(23.0); } @@ -64,7 +71,7 @@ public class TsRollingArgumentEntryTest { values.put(ts - 5, 1.0); newEntry.setTsRecords(values); - assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.updateEntry(newEntry, ctx)).isTrue(); assertThat(entry.getTsRecords()).hasSize(5); assertThat(entry.getTsRecords()).isEqualTo(Map.of( ts - 40, 10.0, @@ -79,7 +86,7 @@ public class TsRollingArgumentEntryTest { void testUpdateEntryWhenValueIsNotNumber() { SingleValueArgumentEntry newEntry = new SingleValueArgumentEntry(ts - 10, new StringDataEntry("key", "string"), 123L); - assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.updateEntry(newEntry, ctx)).isTrue(); assertThat(entry.getTsRecords().get(ts - 10)).isNaN(); } @@ -93,7 +100,7 @@ public class TsRollingArgumentEntryTest { newEntry.setTsRecords(values); entry = new TsRollingArgumentEntry(3, 30000L); - assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.updateEntry(newEntry, ctx)).isTrue(); assertThat(entry.getTsRecords()).hasSize(1); assertThat(entry.getTsRecords()).isEqualTo(Map.of( ts - 5, 0.0 @@ -111,7 +118,7 @@ public class TsRollingArgumentEntryTest { newEntry.setTsRecords(values); entry = new TsRollingArgumentEntry(3, 30000L); - assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.updateEntry(newEntry, ctx)).isTrue(); assertThat(entry.getTsRecords()).hasSize(3); assertThat(entry.getTsRecords()).isEqualTo(Map.of( ts - 18, 0.0,