Browse Source

added test

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
f42c62ac74
  1. 8
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java
  2. 23
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  3. 176
      application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java

8
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;
}
}

23
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<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception {
prepareIntervals();
createIntervalIfNotExist();
prepareIntervals();
long now = System.currentTimeMillis();
Map<AggIntervalEntry, Map<String, ArgumentEntry>> 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<AggIntervalEntry, Map<String, ArgumentEntry>> 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<AggIntervalEntry, Map<String, ArgumentEntry>> 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<AggIntervalEntry, Map<String, ArgumentEntry>> 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);
}

176
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<String, Argument> arguments = new HashMap<>();
Argument argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null));
argument.setLimit(100);
arguments.put("en", argument);
Map<String, AggMetric> 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<String, Argument> inputs,
Map<String, AggMetric> 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);
}
}
Loading…
Cancel
Save