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 e8174967a5..741c94c796 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 @@ -159,6 +159,11 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, } else if (entry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().lastEntry(); newTs = (lastEntry != null) ? lastEntry.getKey() : System.currentTimeMillis(); + } else if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { + newTs = relatedEntitiesArgumentEntry.getEntityInputs().values().stream() + .mapToLong(e -> (e instanceof SingleValueArgumentEntry s) ? s.getTs() : 0L) + .max() + .orElse(0L); } this.latestTimestamp = Math.max(this.latestTimestamp, newTs); } diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesAggregationCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesAggregationCalculatedFieldStateTest.java new file mode 100644 index 0000000000..f1e735fc3b --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesAggregationCalculatedFieldStateTest.java @@ -0,0 +1,234 @@ +/** + * 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.service.cf.ctx.state; + +import com.fasterxml.jackson.databind.JsonNode; +import io.micrometer.core.instrument.simple.SimpleMeterRegistry; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.bean.override.mockito.MockitoBean; +import org.thingsboard.script.api.tbel.DefaultTbelInvokeService; +import org.thingsboard.script.api.tbel.TbelInvokeService; +import org.thingsboard.server.actors.ActorSystemContext; +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.AggFunctionInput; +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.RelatedEntitiesAggregationCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationPathLevel; +import org.thingsboard.server.common.stats.DefaultStatsFactory; +import org.thingsboard.server.dao.usagerecord.ApiLimitService; +import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; +import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAggregationCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesArgumentEntry; + +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.when; + +@SpringBootTest(classes = {SimpleMeterRegistry.class, DefaultStatsFactory.class, DefaultTbelInvokeService.class}) +public class RelatedEntitiesAggregationCalculatedFieldStateTest { + + private final long ts = System.currentTimeMillis(); + + private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("cceba360-71e9-44ab-8596-c600d60ee0d0")); + private final AssetProfileId assetProfileId = new AssetProfileId(UUID.fromString("ccced83b-e5f6-4978-bba0-c2df46de9a35")); + private final AssetId assetId = new AssetId(UUID.fromString("982dfdee-b2bc-4f04-a49d-7ffd70940c69")); + private final DeviceId device1 = new DeviceId(UUID.fromString("47f1fef5-a3b7-46e7-9732-b7669d3ef885")); + private final DeviceId device2 = new DeviceId(UUID.fromString("3b097e5e-9eef-46f4-8997-937075e6f342")); + + private RelatedEntitiesAggregationCalculatedFieldState state; + private CalculatedFieldCtx ctx; + + @Autowired + private TbelInvokeService tbelInvokeService; + + @MockitoBean + private ApiLimitService apiLimitService; + + @MockitoBean + private ActorSystemContext actorSystemContext; + + @BeforeEach + void setUp() { + when(actorSystemContext.getTbelInvokeService()).thenReturn(tbelInvokeService); + when(actorSystemContext.getApiLimitService()).thenReturn(apiLimitService); + when(apiLimitService.getLimit(any(), any())).thenReturn(1000L); + + initCtxAndState(); + } + + void initCtxAndState() { + ctx = new CalculatedFieldCtx(getCalculatedField(), actorSystemContext); + ctx.init(); + + state = new RelatedEntitiesAggregationCalculatedFieldState(assetId); + state.setCtx(ctx, null); + state.init(false); + } + + @Test + void testType() { + assertThat(state.getType()).isEqualTo(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION); + } + + @Test + void testInitAddsRequiredArgument() { + assertThat(state.getRequiredArguments()).contains("oc"); + } + + @Test + void testIsReadyReturnFalseWhenNoArgumentsSet() { + assertThat(state.isReady()).isFalse(); + } + + @Test + void testUpdateArguments() { + assertThat(state.getLastArgsRefreshTs()).isEqualTo(-1); + + state.update(arguments(), ctx); + + assertThat(state.getLastArgsRefreshTs()).isGreaterThan(-1); + } + + @Test + void testPerformCalculationWhenArgsUpdatedButIntervalDidNotPass() throws Exception { + // deduplication interval 60 sec + state.setLastMetricsEvalTs(ts - TimeUnit.SECONDS.toMillis(10)); + state.setLastArgsRefreshTs(ts - TimeUnit.SECONDS.toMillis(5)); + + assertThat(state.performCalculation(arguments(), ctx).get()).isEqualTo(TelemetryCalculatedFieldResult.EMPTY); + } + + @Test + void testPerformCalculationWhenIntervalPassedButNoArgsUpdated() throws Exception { + state.setLastMetricsEvalTs(ts - TimeUnit.SECONDS.toMillis(90)); + state.setLastArgsRefreshTs(ts - TimeUnit.SECONDS.toMillis(90)); + + assertThat(state.performCalculation(arguments(), ctx).get()).isEqualTo(TelemetryCalculatedFieldResult.EMPTY); + } + + @Test + void testPerformCalculationWhenIntervalPassedAndArgsUpdated() throws Exception { + state.setLastMetricsEvalTs(ts - TimeUnit.SECONDS.toMillis(70)); + state.setLastArgsRefreshTs(ts - TimeUnit.SECONDS.toMillis(10)); + + Map arguments = arguments(); + state.update(arguments, ctx); + + TelemetryCalculatedFieldResult calculatedFieldResult = (TelemetryCalculatedFieldResult) state.performCalculation(arguments, ctx).get(); + assertThat(calculatedFieldResult).isNotNull(); + assertThat(calculatedFieldResult.isEmpty()).isFalse(); + assertThat(calculatedFieldResult.getType()).isEqualTo(OutputType.TIME_SERIES); + + JsonNode result = calculatedFieldResult.getResult(); + assertThat(result).isNotNull(); + safeAssert(result.get("ts"), String.valueOf(state.getLatestTimestamp())); + JsonNode values = result.get("values"); + safeAssert(values.get("freeSpaces"), "1"); + safeAssert(values.get("occupiedSpaces"), "1"); + safeAssert(values.get("totalSpaces"), "2"); + } + + private void safeAssert(JsonNode node, String expected) { + assertThat(node).isNotNull(); + assertThat(node.asText()).isEqualTo(expected); + } + + private Map arguments() { + Map arguments = new HashMap<>(); + Map entityArguments = new HashMap<>(); + entityArguments.put(device1, new SingleValueArgumentEntry(ts - 100, new BooleanDataEntry("occupied", true), 23L)); + entityArguments.put(device2, new SingleValueArgumentEntry(ts - 80, new BooleanDataEntry("occupied", false), 26L)); + RelatedEntitiesArgumentEntry entry = new RelatedEntitiesArgumentEntry(entityArguments, false); + arguments.put("oc", entry); + return arguments; + } + + private CalculatedField getCalculatedField() { + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setTenantId(tenantId); + calculatedField.setEntityId(assetProfileId); + calculatedField.setType(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION); + calculatedField.setName("Test CF"); + calculatedField.setConfigurationVersion(1); + + var config = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + config.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)); + + Argument oc = new Argument(); + ReferencedEntityKey refKey = new ReferencedEntityKey("occupied", ArgumentType.TS_LATEST, null); + oc.setRefEntityKey(refKey); + config.setArguments(Map.of("oc", oc)); + + Map aggMetrics = new HashMap<>(); + + AggMetric freeSpaces = new AggMetric(); + freeSpaces.setFunction(AggFunction.COUNT); + freeSpaces.setFilter("return oc == false;"); + freeSpaces.setInput(new AggKeyInput("oc")); + aggMetrics.put("freeSpaces", freeSpaces); + + AggMetric occupiedSpaces = new AggMetric(); + occupiedSpaces.setFunction(AggFunction.COUNT); + occupiedSpaces.setFilter("return oc == true;"); + occupiedSpaces.setInput(new AggKeyInput("oc")); + aggMetrics.put("occupiedSpaces", occupiedSpaces); + + AggMetric totalSpaces = new AggMetric(); + totalSpaces.setFunction(AggFunction.COUNT); + totalSpaces.setInput(new AggFunctionInput("return 1;")); + aggMetrics.put("totalSpaces", totalSpaces); + config.setMetrics(aggMetrics); + + config.setDeduplicationIntervalInSec(60); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(0); + config.setOutput(output); + + config.setUseLatestTs(true); + + calculatedField.setConfiguration(config); + calculatedField.setVersion(1L); + return calculatedField; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java index cf7040c4bb..0dee6ee4a4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java @@ -38,6 +38,8 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A @Valid @NotEmpty private Map metrics; + @Valid + @NotNull private Output output; private boolean useLatestTs; @@ -55,7 +57,22 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A @Override public void validate() { + validateRelation(); + validateArguments(); + validateMetrics(); + } + + private void validateRelation() { + if (relation == null) { + throw new IllegalArgumentException("Relation must be specified!"); + } relation.validate(); + } + + private void validateArguments() { + if (arguments == null || arguments.isEmpty()) { + throw new IllegalArgumentException("Arguments map cannot be empty."); + } if (arguments.containsKey("ctx")) { throw new IllegalArgumentException("Argument name 'ctx' is reserved and cannot be used."); } @@ -64,4 +81,20 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A } } + private void validateMetrics() { + if (metrics == null || metrics.isEmpty()) { + throw new IllegalArgumentException("Metrics map cannot be empty."); + } + + for (AggMetric metric : metrics.values()) { + if (metric.getInput() instanceof AggKeyInput aggKeyInput) { + if (!arguments.containsKey(aggKeyInput.getKey())) { + throw new IllegalArgumentException( + "Metric references unknown argument: '" + aggKeyInput.getKey() + "'." + ); + } + } + } + } + } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfigurationTest.java new file mode 100644 index 0000000000..c06b6e248d --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfigurationTest.java @@ -0,0 +1,123 @@ +/** + * 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; + +import org.junit.jupiter.api.Test; +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.ReferencedEntityKey; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationPathLevel; + +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +public class RelatedEntitiesAggregationCalculatedFieldConfigurationTest { + + @Test + void typeShouldBeEntityAggregation() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + assertThat(cfg.getType()).isEqualTo(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION); + } + + @Test + void validateShouldThrowWhenRelationIsNotSet() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Relation must be specified!"); + } + + @Test + void validateShouldThrowWhenRelationIsNotValid() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + cfg.setRelation(new RelationPathLevel(null, null)); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Direction must be specified!"); + } + + @Test + void validateShouldThrowWhenArgumentsMapIsEmpty() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Arguments map cannot be empty."); + } + + @Test + void validateShouldThrowWhenTsRollingArgumentUsed() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)); + Argument argument = new Argument(); + argument.setRefEntityKey(new ReferencedEntityKey("key", ArgumentType.TS_ROLLING, null)); + cfg.setArguments(Map.of("k", argument)); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Calculated field with type: '" + CalculatedFieldType.RELATED_ENTITIES_AGGREGATION + "' doesn't support TS_ROLLING arguments."); + } + + @Test + void validateShouldThrowWhenMetricMapIsEmpty() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)); + cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); + cfg.setMetrics(Map.of()); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Metrics map cannot be empty."); + } + + @Test + void validateShouldThrowWhenMetricReferencesUnknownArgument() { + var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)); + cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); + + AggMetric metric = new AggMetric(); + metric.setInput(new AggKeyInput("unknown")); + cfg.setMetrics(Map.of("m", metric)); + + cfg.setOutput(new Output()); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Metric references unknown argument: 'unknown'."); + } + + private Argument validArgument(ArgumentType type) { + Argument a = new Argument(); + a.setRefEntityKey(new ReferencedEntityKey("key", type, null)); + return a; + } + +} diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java index e3ad7e55f1..f352171382 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java @@ -26,9 +26,12 @@ import org.apache.http.ssl.SSLContexts; import org.testng.annotations.AfterSuite; import org.testng.annotations.BeforeSuite; import org.testng.annotations.Listeners; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; +import org.thingsboard.server.common.data.EntityInfo; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.device.profile.AllowCreateNewDevicesDeviceProfileProvisionConfiguration; import org.thingsboard.server.common.data.device.profile.CheckPreProvisionedDevicesDeviceProfileProvisionConfiguration; import org.thingsboard.server.common.data.device.profile.DeviceProfileData; @@ -39,6 +42,7 @@ import org.thingsboard.server.common.data.id.DeviceId; import java.net.URI; import java.util.Map; import java.util.Random; +import java.util.function.Consumer; @Slf4j @@ -186,4 +190,12 @@ public abstract class AbstractContainerTest { return testRestClient.postDeviceProfile(deviceProfile); } + protected void updateDefaultTenantProfile(Consumer updater) { + EntityInfo defaultTenantProfileInfo = testRestClient.getDefaultTenantProfileInfo(); + TenantProfile oldTenantProfile = testRestClient.getTenantProfileById(defaultTenantProfileInfo.getId().getId().toString()); + TenantProfile tenantProfile = JacksonUtil.clone(oldTenantProfile); + updater.accept(tenantProfile); + testRestClient.postTenantProfile(tenantProfile); + } + } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java index 7a4d3f1cfe..36f0ff5647 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java @@ -34,10 +34,12 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.EntityInfo; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.asset.Asset; @@ -816,4 +818,33 @@ public class TestRestClient { .then() .statusCode(HTTP_OK); } + + public TenantProfile postTenantProfile(TenantProfile tenantProfile) { + return given().spec(requestSpec).body(tenantProfile) + .post("/api/tenantProfile") + .then() + .statusCode(HTTP_OK) + .extract() + .as(TenantProfile.class); + } + + public EntityInfo getDefaultTenantProfileInfo() { + return given().spec(requestSpec) + .get("/api/tenantProfileInfo/default") + .then() + .statusCode(HTTP_OK) + .extract() + .as(EntityInfo.class); + } + + public TenantProfile getTenantProfileById(String tenantProfileId) { + return given().spec(requestSpec) + .pathParams("tenantProfileId", tenantProfileId) + .get("/api/tenantProfile/{tenantProfileId}") + .then() + .statusCode(HTTP_OK) + .extract() + .as(TenantProfile.class); + } + } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java index 9c043eee8d..b398ef73f4 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java @@ -38,6 +38,11 @@ import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunctionInput; +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.RelatedEntitiesAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.ZoneGroupConfiguration; @@ -53,6 +58,8 @@ import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationPathLevel; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.msa.AbstractContainerTest; import org.thingsboard.server.msa.ui.utils.EntityPrototypes; @@ -100,6 +107,14 @@ public class CalculatedFieldTest extends AbstractContainerTest { public void beforeClass() { testRestClient.login("sysadmin@thingsboard.org", "sysadmin"); + updateDefaultTenantProfile(tenantProfile -> { + TenantProfileData profileData = tenantProfile.getProfileData(); + DefaultTenantProfileConfiguration profileConfiguration = (DefaultTenantProfileConfiguration) profileData.getConfiguration(); + profileConfiguration.setMinAllowedDeduplicationIntervalInSecForCF(1); + profileConfiguration.setMinAllowedScheduledUpdateIntervalInSecForCF(1); + tenantProfile.setProfileData(profileData); + }); + tenantId = testRestClient.postTenant(EntityPrototypes.defaultTenantPrototype("Tenant")).getId(); tenantAdminId = testRestClient.createUserAndLogin(defaultTenantAdmin(tenantId, "tenantAdmin@thingsboard.org"), "tenant"); @@ -592,6 +607,187 @@ public class CalculatedFieldTest extends AbstractContainerTest { testRestClient.deleteCalculatedFieldIfExists(saved.getId()); } + @Test + public void testRelatedEntitiesAggregationCalculatedField() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + // --- Create entities --- + String device_1_1_token = "000000011"; + Device device_1_1 = testRestClient.postDevice(device_1_1_token, createDevice("Device 1-1", deviceProfileId)); + String device_1_2_token = "000000012"; + Device device_1_2 = testRestClient.postDevice(device_1_2_token, createDevice("Device 1-2", deviceProfileId)); + + // Create relations FROM asset TO devices + EntityRelation rel_1_1 = new EntityRelation(asset.getId(), device_1_1.getId(), EntityRelation.CONTAINS_TYPE); + EntityRelation rel_1_2 = new EntityRelation(asset.getId(), device_1_2.getId(), EntityRelation.CONTAINS_TYPE); + testRestClient.postEntityRelation(rel_1_1); + testRestClient.postEntityRelation(rel_1_2); + + // Post telemetry + testRestClient.postTelemetry(device_1_1_token, JacksonUtil.toJsonNode("{\"occupied\":true}")); + testRestClient.postTelemetry(device_1_2_token, JacksonUtil.toJsonNode("{\"occupied\":false}")); + + // --- Create CF: Related entities aggregation --- + CalculatedField calculatedField = createOccupancyCF(assetProfileId); + + // --- Assert aggregation --- + await().alias("create cf -> check aggregation") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode occupancy = testRestClient.getLatestTelemetry(asset.getId()); + assertThat(occupancy).isNotNull(); + + assertThat(occupancy.get("freeSpaces")).isNotNull(); + assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1"); + + assertThat(occupancy.get("occupiedSpaces")).isNotNull(); + assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("1"); + + assertThat(occupancy.get("totalSpaces")).isNotNull(); + assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2"); + }); + + // Post telemetry + testRestClient.postTelemetry(device_1_2_token, JacksonUtil.toJsonNode("{\"occupied\":true}")); + + // --- Assert aggregation --- + await().alias("update telemetry -> check aggregation") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode occupancy = testRestClient.getLatestTelemetry(asset.getId()); + assertThat(occupancy).isNotNull(); + + assertThat(occupancy.get("freeSpaces")).isNotNull(); + assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("0"); + + assertThat(occupancy.get("occupiedSpaces")).isNotNull(); + assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("2"); + + assertThat(occupancy.get("totalSpaces")).isNotNull(); + assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2"); + }); + + // Add entity to profile + Asset asset2 = testRestClient.postAsset(createAsset("Asset 2", assetProfileId)); + String device_2_1_token = "000000021"; + Device device_2_1 = testRestClient.postDevice(device_2_1_token, createDevice("Device 2-1", deviceProfileId)); + String device_2_2_token = "000000022"; + Device device_2_2 = testRestClient.postDevice(device_2_2_token, createDevice("Device 2-2", deviceProfileId)); + + // Post telemetry + testRestClient.postTelemetry(device_2_1_token, JacksonUtil.toJsonNode("{\"occupied\":true}")); + testRestClient.postTelemetry(device_2_2_token, JacksonUtil.toJsonNode("{\"occupied\":false}")); + + // --- Assert aggregation --- + await().alias("add entity to profile cf -> no aggregated values since no relations") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode occupancy = testRestClient.getLatestTelemetry(asset2.getId()); + assertThat(occupancy).isNullOrEmpty(); + }); + + // Create relations FROM asset TO devices + EntityRelation rel_2_1 = new EntityRelation(asset2.getId(), device_2_1.getId(), EntityRelation.CONTAINS_TYPE); + testRestClient.postEntityRelation(rel_2_1); + EntityRelation rel_2_2 = new EntityRelation(asset2.getId(), device_2_2.getId(), EntityRelation.CONTAINS_TYPE); + testRestClient.postEntityRelation(rel_2_2); + + // --- Assert aggregation --- + await().alias("create relation -> check aggregation") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode occupancy = testRestClient.getLatestTelemetry(asset2.getId()); + assertThat(occupancy).isNotNull(); + + assertThat(occupancy.get("freeSpaces")).isNotNull(); + assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("1"); + + assertThat(occupancy.get("occupiedSpaces")).isNotNull(); + assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("1"); + + assertThat(occupancy.get("totalSpaces")).isNotNull(); + assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("2"); + }); + + testRestClient.deleteEntityRelation(asset2.getId(), EntityRelation.CONTAINS_TYPE, device_2_2.getId()); + + // --- Assert aggregation --- + await().alias("delete relation -> check aggregation") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode occupancy = testRestClient.getLatestTelemetry(asset2.getId()); + assertThat(occupancy).isNotNull(); + + assertThat(occupancy.get("freeSpaces")).isNotNull(); + assertThat(occupancy.get("freeSpaces").get(0).get("value").asText()).isEqualTo("0"); + + assertThat(occupancy.get("occupiedSpaces")).isNotNull(); + assertThat(occupancy.get("occupiedSpaces").get(0).get("value").asText()).isEqualTo("1"); + + assertThat(occupancy.get("totalSpaces")).isNotNull(); + assertThat(occupancy.get("totalSpaces").get(0).get("value").asText()).isEqualTo("1"); + }); + + testRestClient.deleteCalculatedFieldIfExists(calculatedField.getId()); + } + + private CalculatedField createOccupancyCF(EntityId entityId) { + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setName("Occupancy"); + calculatedField.setEntityId(entityId); + calculatedField.setType(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION); + + RelatedEntitiesAggregationCalculatedFieldConfiguration configuration = new RelatedEntitiesAggregationCalculatedFieldConfiguration(); + + configuration.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, "Contains")); + + Map arguments = new HashMap<>(); + Argument argument = new Argument(); + argument.setRefEntityKey(new ReferencedEntityKey("occupied", ArgumentType.TS_LATEST, null)); + argument.setDefaultValue("false"); + arguments.put("oc", argument); + configuration.setArguments(arguments); + + configuration.setDeduplicationIntervalInSec(5); + configuration.setScheduledUpdateInterval(10); + + Map aggMetrics = new HashMap<>(); + + AggMetric freeSpaces = new AggMetric(); + freeSpaces.setFunction(AggFunction.COUNT); + freeSpaces.setFilter("return oc == false;"); + freeSpaces.setInput(new AggKeyInput("oc")); + aggMetrics.put("freeSpaces", freeSpaces); + + AggMetric occupiedSpaces = new AggMetric(); + occupiedSpaces.setFunction(AggFunction.COUNT); + occupiedSpaces.setFilter("return oc == true;"); + occupiedSpaces.setInput(new AggKeyInput("oc")); + aggMetrics.put("occupiedSpaces", occupiedSpaces); + + AggMetric totalSpaces = new AggMetric(); + totalSpaces.setFunction(AggFunction.COUNT); + totalSpaces.setInput(new AggFunctionInput("return 1;")); + aggMetrics.put("totalSpaces", totalSpaces); + configuration.setMetrics(aggMetrics); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(0); + configuration.setOutput(output); + + calculatedField.setConfiguration(configuration); + calculatedField.setDebugSettings(DebugSettings.all()); + + return testRestClient.postCalculatedField(calculatedField); + } + private CalculatedField createSimpleCalculatedField() { return createSimpleCalculatedField(device.getId()); }