From f42c62ac7481a80cfda9acdc5c453d7f62391f15 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 3 Nov 2025 09:19:17 +0200 Subject: [PATCH] added test --- .../single/AggIntervalEntryStatus.java | 8 + ...EntityAggregationCalculatedFieldState.java | 23 ++- .../EntityAggregationCalculatedFieldTest.java | 176 ++++++++++++++++++ 3 files changed, 200 insertions(+), 7 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java index fa9bbd5a54..fd2a899503 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java @@ -45,4 +45,12 @@ public class AggIntervalEntryStatus { return false; } + public boolean intervalPassed(long checkInterval) { + boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - checkInterval; + if (intervalPassed) { + lastMetricsEvalTs = System.currentTimeMillis(); + } + return intervalPassed; + } + } 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 6885aa8b42..5510d63f24 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 @@ -41,6 +41,8 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; + public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { private AggInterval interval; @@ -58,8 +60,8 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt } public void scheduleReevaluation() { - fillMissingIntervals(interval.getCurrentIntervalEndTs(), intervalDuration); prepareIntervals(); + fillMissingIntervals(interval.getCurrentIntervalEndTs(), intervalDuration); long now = System.currentTimeMillis(); intervals.forEach((intervalEntry, argumentIntervalStatuses) -> { if (intervalEntry.belongsToInterval(now)) { @@ -115,8 +117,8 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt @Override public ListenableFuture performCalculation(Map updatedArgs, CalculatedFieldCtx ctx) throws Exception { - prepareIntervals(); createIntervalIfNotExist(); + prepareIntervals(); long now = System.currentTimeMillis(); Map> results = new HashMap<>(); @@ -163,8 +165,10 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt } arguments.forEach((argName, argumentEntry) -> { var entityAggEntry = (EntityAggregationArgumentEntry) argumentEntry; - entityAggEntry.getAggIntervals().put(currentInterval, new AggIntervalEntryStatus()); - intervals.computeIfAbsent(currentInterval, i -> new HashMap<>()).put(argName, new AggIntervalEntryStatus()); + if (!entityAggEntry.getAggIntervals().containsKey(currentInterval)) { + entityAggEntry.getAggIntervals().put(currentInterval, new AggIntervalEntryStatus()); + intervals.computeIfAbsent(currentInterval, i -> new HashMap<>()).put(argName, new AggIntervalEntryStatus()); + } }); ctx.scheduleReevaluation(interval.getDelayUntilIntervalEnd(), actorCtx); } @@ -190,7 +194,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results) { args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { - processMetric(intervalEntry, argName, results); + processMetric(intervalEntry, argName, false, results); } }); } @@ -200,20 +204,25 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results) { args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.shouldRecalculate(checkInterval)) { - processMetric(intervalEntry, argName, results); + processMetric(intervalEntry, argName, false, results); ctx.scheduleReevaluation(checkInterval, actorCtx); + } else if (argEntryIntervalStatus.intervalPassed(checkInterval)) { + processMetric(intervalEntry, argName, true, results); } }); } private void processMetric(AggIntervalEntry intervalEntry, String argName, + boolean useDefault, Map> results) { String metricName = findMetricName(argName); if (metricName != null) { AggMetric metric = metrics.get(metricName); String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); - ArgumentEntry metricEntry = cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); + ArgumentEntry metricEntry = useDefault + ? ArgumentEntry.createSingleValueArgument(createDefaultKvEntry(argKey, metric.getDefaultValue())) + : cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); if (!metricEntry.isEmpty()) { results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry); } diff --git a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java new file mode 100644 index 0000000000..2509263dce --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -0,0 +1,176 @@ +/** + * 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.cf; + +import com.fasterxml.jackson.databind.node.ObjectNode; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.springframework.test.annotation.DirtiesContext; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; +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.AggIntervalType; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.CustomInterval; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; +import org.thingsboard.server.common.data.debug.DebugSettings; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.controller.AbstractControllerTest; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; +import static org.thingsboard.server.cf.CalculatedFieldIntegrationTest.POLL_INTERVAL; + +@DaoSqlTest +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD) +public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest { + + private Tenant savedTenant; + + @Before + public void beforeEach() throws Exception { + loginSysAdmin(); + + updateDefaultTenantProfileConfig(tenantProfileConfig -> { + tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1); + }); + + Tenant tenant = new Tenant(); + tenant.setTitle("My tenant"); + savedTenant = saveTenant(tenant); + assertThat(savedTenant).isNotNull(); + + User tenantAdmin = new User(); + tenantAdmin.setAuthority(Authority.TENANT_ADMIN); + tenantAdmin.setTenantId(savedTenant.getId()); + tenantAdmin.setEmail("tenant@thingsboard.org"); + tenantAdmin.setFirstName("John"); + tenantAdmin.setLastName("Doe"); + + createUserAndLogin(tenantAdmin, "testPassword"); + } + + @After + public void afterTest() throws Exception { + loginSysAdmin(); + + deleteTenant(savedTenant.getId()); + } + + @Test + public void testCreateCf_checkAggregation() throws Exception { + Device device = createDevice("Device", "1234567890111"); + + CustomInterval customInterval = new CustomInterval(1, AggIntervalType.MIN, 0, "Europe/Kyiv"); + long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); + long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); + + long tsBeforeInterval = currentIntervalStartTs - 1000L; + long tsInInterval_1 = currentIntervalStartTs + 1000L; + long tsInInterval_2 = currentIntervalStartTs + 500L; + long tsInInterval_3 = currentIntervalStartTs + 200L; + 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.getIntervalDurationMillis(); + CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval); + + await().alias("create CF and perform aggregation after interval end") + .atMost(2 * interval, TimeUnit.MILLISECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode result = getLatestTelemetry(device.getId(), "consumptionPerMin"); + assertThat(result).isNotNull(); + assertThat(result.get("consumptionPerMin").get(0).get("value").asText()).isEqualTo("400"); + }); + } + + private CalculatedField createTotalConsumptionCF(EntityId entityId, AggInterval aggInterval) { + Map arguments = new HashMap<>(); + Argument argument = new Argument(); + argument.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null)); + argument.setLimit(100); + arguments.put("en", argument); + + Map aggMetrics = new HashMap<>(); + + AggMetric consumptionPerMin = new AggMetric(); + consumptionPerMin.setFunction(AggFunction.SUM); + consumptionPerMin.setInput(new AggKeyInput("en")); + aggMetrics.put("consumptionPerMin", consumptionPerMin); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(0); + + return createAggCf("Consumption per minute", entityId, + aggInterval, + new Watermark(TimeUnit.MINUTES.toMillis(1), TimeUnit.SECONDS.toMillis(10)), + arguments, + aggMetrics, + output); + } + + private CalculatedField createAggCf(String name, + EntityId entityId, + AggInterval aggInterval, + Watermark watermark, + Map inputs, + Map metrics, + Output output) { + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setName(name); + calculatedField.setEntityId(entityId); + calculatedField.setType(CalculatedFieldType.ENTITY_AGGREGATION); + + EntityAggregationCalculatedFieldConfiguration configuration = new EntityAggregationCalculatedFieldConfiguration(); + + configuration.setArguments(inputs); + configuration.setMetrics(metrics); + configuration.setInterval(aggInterval); + configuration.setWatermark(watermark); + configuration.setOutput(output); + + calculatedField.setConfiguration(configuration); + calculatedField.setDebugSettings(DebugSettings.all()); + return saveCalculatedField(calculatedField); + } + + private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception { + return doGetAsync("/api/plugins/telemetry/" + entityId.getEntityType() + "/" + entityId.getId() + "/values/timeseries?keys=" + String.join(",", keys), ObjectNode.class); + } + +}