diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 4e6711b150..e6dc7cbcd3 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -262,17 +262,35 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } private void onTenantProfileUpdated(ComponentLifecycleMsg msg, TbCallback callback) { + checkCfIntervalForUpdate(); + + long maxRelatedEntitiesPerCfArgument = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxRelatedEntitiesToReturnPerCfArgument); + List cfsToReinit = new ArrayList<>(); + Stream.concat( + calculatedFields.values().stream(), + entityIdCalculatedFields.values().stream().flatMap(Collection::stream) + ).forEach(ctx -> { + if (ctx.hasRelatedEntities() && ctx.getMaxRelatedEntitiesPerCfArgument() != maxRelatedEntitiesPerCfArgument) { + cfsToReinit.add(ctx); + } + ctx.setTenantProfileProperties(); + }); + + if (!cfsToReinit.isEmpty()) { + MultipleTbCallback cfsReinitCallback = new MultipleTbCallback(cfsToReinit.size(), callback); + cfsToReinit.forEach(ctx -> applyToTargetCfEntityActors(ctx, cfsReinitCallback, (id, cb) -> initCfForEntity(id, ctx, StateAction.REINIT, cb))); + } else { + callback.onSuccess(); + } + } + + private void checkCfIntervalForUpdate() { long updatedCfCheckInterval = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getCfReevaluationCheckInterval); if (cfCheckInterval != updatedCfCheckInterval) { cfCheckInterval = updatedCfCheckInterval; cancelReevaluationTask(); scheduleCfsReevaluation(); } - Stream.concat( - calculatedFields.values().stream(), - entityIdCalculatedFields.values().stream().flatMap(Collection::stream) - ).forEach(CalculatedFieldCtx::setTenantProfileProperties); - callback.onSuccess(); } private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) { @@ -383,10 +401,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware .toList(); MultipleTbCallback directionCallback = new MultipleTbCallback(matchingCfs.size(), parentCallback); - - matchingCfs.forEach(ctx -> - applyToTargetCfEntityActors(ctx, directionCallback, (entityId, cb) -> relationAction.accept(entityId, ctx, cb)) - ); + matchingCfs.forEach(ctx -> relationAction.accept(mainId, ctx, directionCallback)); } private void onCfCreated(ComponentLifecycleMsg msg, TbCallback callback) throws CalculatedFieldException { @@ -580,18 +595,12 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware Predicate matchesCfEntity = relatedEntity -> cfEntityId.equals(relatedEntity) || cfEntityId.equals(getProfileId(tenantId, relatedEntity)); if (byRelationPathQuery != null && !byRelationPathQuery.isEmpty()) { switch (relation.direction()) { - case FROM -> { - EntityRelation entityRelation = byRelationPathQuery.get(0); // only one supported - EntityId relatedId = entityRelation.getFrom(); - if (matchesCfEntity.test(relatedId)) { - result.add(new CalculatedFieldEntityCtxId(tenantId, cf.getCfId(), relatedId)); - } - } - case TO -> { - byRelationPathQuery.stream() - .filter(entityRelation -> matchesCfEntity.test(entityRelation.getTo())) - .forEach(entityRelation -> result.add(new CalculatedFieldEntityCtxId(tenantId, cf.getCfId(), entityRelation.getTo()))); - } + case FROM -> byRelationPathQuery.stream() + .filter(entityRelation -> matchesCfEntity.test(entityRelation.getFrom())) + .forEach(entityRelation -> result.add(new CalculatedFieldEntityCtxId(tenantId, cf.getCfId(), entityRelation.getFrom()))); + case TO -> byRelationPathQuery.stream() + .filter(entityRelation -> matchesCfEntity.test(entityRelation.getTo())) + .forEach(entityRelation -> result.add(new CalculatedFieldEntityCtxId(tenantId, cf.getCfId(), entityRelation.getTo()))); } } } diff --git a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java index bca55e9784..ee2e871d6e 100644 --- a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java +++ b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java @@ -164,6 +164,7 @@ public class SystemInfoController extends BaseController { systemParams.setMaxDataPointsPerRollingArg(tenantProfileConfiguration.getMaxDataPointsPerRollingArg()); systemParams.setMinAllowedScheduledUpdateIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedScheduledUpdateIntervalInSecForCF()); systemParams.setMaxRelationLevelPerCfArgument(tenantProfileConfiguration.getMaxRelationLevelPerCfArgument()); + systemParams.setMaxRelatedEntitiesToReturnPerCfArgument(tenantProfileConfiguration.getMaxRelatedEntitiesToReturnPerCfArgument()); systemParams.setMinAllowedDeduplicationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedDeduplicationIntervalInSecForCF()); systemParams.setMinAllowedAggregationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedAggregationIntervalInSecForCF()); systemParams.setIntermediateAggregationIntervalInSecForCF(tenantProfileConfiguration.getIntermediateAggregationIntervalInSecForCF()); 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 fed1e82e96..8cc0043198 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 @@ -247,14 +247,8 @@ public abstract class AbstractCalculatedFieldProcessingService { } return switch (relation.direction()) { - case FROM -> relations.stream() - .map(EntityRelation::getTo) - .toList(); - case TO -> relations.stream() - .map(EntityRelation::getFrom) - .findFirst() - .map(List::of) - .orElseGet(Collections::emptyList); + case FROM -> relations.stream().map(EntityRelation::getTo).toList(); + case TO -> relations.stream().map(EntityRelation::getFrom).toList(); }; }, calculatedFieldCallbackExecutor); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index 6500bfb892..e5215eca28 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -29,6 +29,7 @@ import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActorRef; import org.thingsboard.server.actors.calculatedField.CalculatedFieldReevaluateMsg; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.alarm.rule.AlarmRule; import org.thingsboard.server.common.data.alarm.rule.condition.expression.TbelAlarmConditionExpression; import org.thingsboard.server.common.data.cf.CalculatedField; @@ -54,11 +55,9 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BasicKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; -import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.util.CollectionsUtil; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.relation.RelationService; -import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.dao.util.TimeUtils; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; @@ -129,6 +128,7 @@ public class CalculatedFieldCtx implements Closeable { private long scheduledUpdateIntervalMillis; private long cfCheckReevaluationIntervalMillis; private long alarmReevaluationIntervalMillis; + private long maxRelatedEntitiesPerCfArgument; private Argument propagationArgument; private boolean applyExpressionForResolvedArguments; @@ -301,12 +301,18 @@ public class CalculatedFieldCtx implements Closeable { } public void setTenantProfileProperties() { - ApiLimitService apiLimitService = systemContext.getApiLimitService(); - this.maxStateSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; - this.maxSingleValueArgumentSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; - this.intermediateAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getIntermediateAggregationIntervalInSecForCF)); - this.cfCheckReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getCfReevaluationCheckInterval)); - this.alarmReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getAlarmsReevaluationInterval)); + TenantProfile tenantProfile = systemContext.getTenantProfileCache().get(tenantId); + if (tenantProfile == null) { + throw new IllegalStateException("Tenant Profile not found for tenant: " + tenantId); + } + tenantProfile.getProfileConfiguration().ifPresent(config -> { + this.maxStateSize = config.getMaxStateSizeInKBytes() * 1024L; + this.maxSingleValueArgumentSize = config.getMaxSingleValueArgumentSizeInKBytes() * 1024L; + this.intermediateAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(config.getIntermediateAggregationIntervalInSecForCF()); + this.cfCheckReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(config.getCfReevaluationCheckInterval()); + this.alarmReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(config.getAlarmsReevaluationInterval()); + this.maxRelatedEntitiesPerCfArgument = config.getMaxRelatedEntitiesToReturnPerCfArgument(); + }); } public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) { @@ -756,6 +762,12 @@ public class CalculatedFieldCtx implements Closeable { return scheduledUpdateIntervalMillis == DISABLED_INTERVAL_VALUE; } + public boolean hasRelatedEntities() { + return CalculatedFieldType.GEOFENCING == cfType + || CalculatedFieldType.PROPAGATION == cfType + || CalculatedFieldType.RELATED_ENTITIES_AGGREGATION == cfType; + } + public boolean shouldFetchRelatedEntities(CalculatedFieldState state) { if (!cfHasRelationPathQuerySource) { return false; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java index 357e3b66d3..cb64e03d90 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInp 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.EntityId; +import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.cf.CalculatedFieldResult; import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; @@ -300,8 +301,25 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat if (argumentEntry == null || argumentEntry.isEmpty()) { return ReadinessStatus.notReady(MISSING_AGGREGATION_ENTITIES_ERROR); } + if (argumentEntry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { + try { + checkConstraintByDirection(relatedEntitiesArgumentEntry); + } catch (Exception e) { + return ReadinessStatus.notReady(e.getMessage()); + } + } } return ReadinessStatus.READY; } + public void checkConstraintByDirection(RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { + if (ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) { + if (EntitySearchDirection.TO == config.getRelation().direction()) { + if (relatedEntitiesArgumentEntry.getEntityInputs().size() > 1) { + throw new IllegalArgumentException("More than one related entity is not supported for relation direction 'TO'. Found: " + relatedEntitiesArgumentEntry.getEntityInputs().size() + "."); + } + } + } + } + } 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 1939318e7a..7392ff1e0f 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 @@ -66,23 +66,34 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs @Override public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { + checkRelatedEntitiesNumber(ctx); entityInputs.putAll(relatedEntitiesArgumentEntry.entityInputs); - return true; } else if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { if (entry.isForceResetPrevious()) { + checkRelatedEntitiesNumber(ctx); entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); - return true; - } - ArgumentEntry argumentEntry = entityInputs.get(singleValueArgumentEntry.getEntityId()); - if (argumentEntry != null) { - argumentEntry.updateEntry(singleValueArgumentEntry, ctx); } else { - entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); + ArgumentEntry argumentEntry = entityInputs.get(singleValueArgumentEntry.getEntityId()); + if (argumentEntry != null) { + argumentEntry.updateEntry(singleValueArgumentEntry, ctx); + } else { + checkRelatedEntitiesNumber(ctx); + entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); + } } - return true; } else { throw new IllegalArgumentException("Unsupported argument entry type for aggregation argument entry: " + entry.getType()); } + return true; + } + + private void checkRelatedEntitiesNumber(CalculatedFieldCtx ctx) { + if (entityInputs.size() >= ctx.getMaxRelatedEntitiesPerCfArgument()) { + throw new IllegalArgumentException( + "Exceeded the maximum allowed related entities per argument '" + + ctx.getMaxRelatedEntitiesPerCfArgument() + "'. Increase the limit in the tenant profile configuration." + ); + } } @Override 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 7b8b393371..c3b5951da6 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 @@ -65,7 +65,7 @@ public class PropagationArgumentEntry implements ArgumentEntry { throw new IllegalArgumentException("Unsupported argument entry type for propagation argument entry: " + entry.getType()); } if (updated.getAdded() != null) { - return checkAdded(updated.getAdded()); + return checkAdded(updated.getAdded(), ctx); } if (updated.getRemoved() != null) { return entityIds.remove(updated.getRemoved()); @@ -80,7 +80,7 @@ public class PropagationArgumentEntry implements ArgumentEntry { return true; } boolean retained = entityIds.retainAll(dbEntityIds); - boolean added = checkAdded(dbEntityIds); + boolean added = checkAdded(dbEntityIds, ctx); return retained || added; } if (updated.isEmpty()) { @@ -91,8 +91,14 @@ public class PropagationArgumentEntry implements ArgumentEntry { return true; } - private boolean checkAdded(Collection updatedIds) { + private boolean checkAdded(Collection updatedIds, CalculatedFieldCtx ctx) { for (EntityId id : updatedIds) { + if (entityIds.size() >= ctx.getMaxRelatedEntitiesPerCfArgument()) { + throw new IllegalArgumentException( + "Exceeded the maximum allowed related entities per argument '" + + ctx.getMaxRelatedEntitiesPerCfArgument() + "'. Increase the limit in the tenant profile configuration." + ); + } if (entityIds.add(id)) { if (added == null) { added = new ArrayList<>(); 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 d23e2900da..fcf8c0593b 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -68,9 +68,9 @@ public class CalculatedFieldArgumentUtils { } public static ArgumentEntry createDefaultMetricArgumentEntry(String argKey, AggMetric metric) { - Long defaultValue = metric.getDefaultValue(); + Double defaultValue = metric.getDefaultValue(); if (defaultValue != null) { - return ArgumentEntry.createSingleValueArgument(new DoubleDataEntry(argKey, defaultValue.doubleValue())); + return ArgumentEntry.createSingleValueArgument(new DoubleDataEntry(argKey, defaultValue)); } return new SingleValueArgumentEntry(); } 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 6c08fc1458..0872e21dae 100644 --- a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -254,7 +254,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest AggMetric consumption = new AggMetric(); consumption.setFunction(AggFunction.SUM); consumption.setInput(new AggKeyInput("en")); - consumption.setDefaultValue(9999L); + consumption.setDefaultValue(9999.0); aggMetrics.put("consumption", consumption); AggMetric avgEnergyConsumption = new AggMetric(); @@ -319,7 +319,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest AggMetric consumption = new AggMetric(); consumption.setFunction(AggFunction.SUM); consumption.setInput(new AggKeyInput("en")); - consumption.setDefaultValue(9999L); + consumption.setDefaultValue(9999.0); aggMetrics.put("consumption", consumption); AggMetric avgTemperature = new AggMetric(); diff --git a/application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java b/application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java index 10a7b46282..4c708bce56 100644 --- a/application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java @@ -470,21 +470,47 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr @Test public void testCreateRelation_checkAggregation() throws Exception { - createOccupancyCF(asset.getId()); - checkInitialCalculation(); + Asset asset2 = createAsset("Asset 2", assetProfile.getId()); + Device device3 = createDevice("Device 3", "1234567890333"); + Device device4 = createDevice("Device 4", "1234567890444"); - Device device3 = createDevice("Device 3", deviceProfile.getId(), "1234567890333"); + createEntityRelation(asset2.getId(), device3.getId(), "Contains"); + createEntityRelation(asset2.getId(), device4.getId(), "Contains"); - postTelemetry(device3.getId(), "{\"occupied\":true}"); + createOccupancyCF(assetProfile.getId()); - createEntityRelation(asset.getId(), device3.getId(), "Contains"); + await().alias("create CF and perform initial aggregation").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + verifyTelemetry(asset.getId(), Map.of( + "freeSpaces", "1", + "occupiedSpaces", "1", + "totalSpaces", "2" + )); - await().alias("create relation and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS) + verifyTelemetry(asset2.getId(), Map.of( + "freeSpaces", "2", + "occupiedSpaces", "0", + "totalSpaces", "2" + )); + }); + + Device device5 = createDevice("Device 5", "1234567890555"); + createEntityRelation(asset2.getId(), device5.getId(), "Contains"); + + await().alias("create relation and perform aggregation on asset 2") + .atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { verifyTelemetry(asset.getId(), Map.of( "freeSpaces", "1", - "occupiedSpaces", "2", + "occupiedSpaces", "1", + "totalSpaces", "2" + )); + + verifyTelemetry(asset2.getId(), Map.of( + "freeSpaces", "3", + "occupiedSpaces", "0", "totalSpaces", "3" )); }); @@ -721,6 +747,43 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr }); } + @Test + public void testUpdateMaxRelatedEntitiesPerArgument_checkAggregation() throws Exception { + loginSysAdmin(); + + updateDefaultTenantProfileConfig(tenantProfileConfig -> { + tenantProfileConfig.setMaxRelatedEntitiesToReturnPerCfArgument(1); + }); + + login("tenant@thingsboard.org", "testPassword"); + + createCountCF(asset.getId()); + + await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode numberOfDevices = getLatestTelemetry(asset.getId(), "numberOfDevices"); + assertThat(numberOfDevices).isNotNull(); + assertThat(numberOfDevices.get("numberOfDevices").get(0).get("value").asText()).isEqualTo("1"); + }); + + loginSysAdmin(); + + updateDefaultTenantProfileConfig(tenantProfileConfig -> { + tenantProfileConfig.setMaxRelatedEntitiesToReturnPerCfArgument(10); + }); + + login("tenant@thingsboard.org", "testPassword"); + + await().alias("update max related entities per argument and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode numberOfDevices = getLatestTelemetry(asset.getId(), "numberOfDevices"); + assertThat(numberOfDevices).isNotNull(); + assertThat(numberOfDevices.get("numberOfDevices").get(0).get("value").asText()).isEqualTo("2"); + }); + } + private void checkInitialCalculation() { await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) @@ -832,6 +895,30 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr output); } + private CalculatedField createCountCF(EntityId entityId) { + Map arguments = new HashMap<>(); + Argument argument = new Argument(); + argument.setRefEntityKey(new ReferencedEntityKey("active", ArgumentType.TS_LATEST, null)); + argument.setDefaultValue("true"); + arguments.put("active", argument); + + Map aggMetrics = new HashMap<>(); + + AggMetric avgMetric = new AggMetric(); + avgMetric.setFunction(AggFunction.COUNT); + avgMetric.setInput(new AggKeyInput("active")); + aggMetrics.put("numberOfDevices", avgMetric); + + TimeSeriesOutput output = new TimeSeriesOutput(); + output.setDecimalsByDefault(0); + + return createAggCf("Number of devices", entityId, + new RelationPathLevel(EntitySearchDirection.FROM, "Contains"), + arguments, + aggMetrics, + output); + } + private CalculatedField createAggCf(String name, EntityId entityId, RelationPathLevel relation, diff --git a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java index 396c381255..d9ec187e39 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -382,7 +382,7 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { AggMetric metric = new AggMetric(); metric.setInput(new AggKeyInput("en")); - metric.setDefaultValue(9999L); + metric.setDefaultValue(9999.0); config.setMetrics(Map.of("consumption", metric)); config.setWatermark(new Watermark(TimeUnit.DAYS.toSeconds(1))); 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 596720d213..7db8bc658a 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 @@ -33,6 +33,7 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.lenient; @ExtendWith(MockitoExtension.class) public class PropagationArgumentEntryTest { @@ -48,6 +49,8 @@ public class PropagationArgumentEntryTest { @BeforeEach void setUp() { + lenient().when(ctx.getMaxRelatedEntitiesPerCfArgument()).thenReturn(1000L); + List propagationEntityIds = new ArrayList<>(); propagationEntityIds.add(ENTITY_1_ID); propagationEntityIds.add(ENTITY_2_ID); 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 db8ce32df6..725860dd80 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 @@ -32,6 +32,7 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.lenient; @ExtendWith(MockitoExtension.class) public class RelatedEntitiesArgumentEntryTest { @@ -48,6 +49,8 @@ public class RelatedEntitiesArgumentEntryTest { @BeforeEach void setUp() { + lenient().when(ctx.getMaxRelatedEntitiesPerCfArgument()).thenReturn(1000L); + Map aggInputs = new HashMap<>(); aggInputs.put(device1, new SingleValueArgumentEntry(device1, new BasicTsKvEntry(ts - 100, new LongDataEntry("key", 12L), 1L))); aggInputs.put(device2, new SingleValueArgumentEntry(device2, new BasicTsKvEntry(ts - 150, new LongDataEntry("key", 16L), 6L))); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java index 678fdc0900..3cd2ada38f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java @@ -40,6 +40,7 @@ public class SystemParams { long maxDataPointsPerRollingArg; int minAllowedScheduledUpdateIntervalInSecForCF; int maxRelationLevelPerCfArgument; + int maxRelatedEntitiesToReturnPerCfArgument; long minAllowedDeduplicationIntervalInSecForCF; long minAllowedAggregationIntervalInSecForCF; long intermediateAggregationIntervalInSecForCF; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java index ea841e24eb..8a322bd9ce 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java @@ -27,6 +27,6 @@ public class AggMetric { private AggFunction function; private String filter; private AggInput input; - private Long defaultValue; + private Double defaultValue; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java index 0fd914df7c..6f404e1aca 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java @@ -52,14 +52,15 @@ import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.eventsourcing.RelationActionEvent; -import org.thingsboard.server.exception.DataValidationException; import org.thingsboard.server.dao.service.ConstraintValidator; import org.thingsboard.server.dao.sql.JpaExecutorService; import org.thingsboard.server.dao.sql.relation.JpaRelationQueryExecutorService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; +import org.thingsboard.server.exception.DataValidationException; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -518,7 +519,13 @@ class BaseRelationService implements RelationService { return Collections.emptyList(); } List relations = relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations; - return relations.size() > limit ? relations.subList(0, limit) : relations; + if (relations.size() > limit) { + List limitedRelations = new ArrayList<>(relations); + limitedRelations.sort(Comparator.comparing(r -> r.getFrom().getId())); + return limitedRelations.subList(0, limit); + } else { + return relations; + } }, directExecutor()); } return executor.submit(() -> { @@ -545,7 +552,13 @@ class BaseRelationService implements RelationService { case FROM -> findByFromAndType(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON); case TO -> findByToAndType(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON); }; - return relations.size() > limit ? relations.subList(0, limit) : relations; + if (relations.size() > limit) { + List limitedRelations = new ArrayList<>(relations); + limitedRelations.sort(Comparator.comparing(r -> r.getFrom().getId())); + return limitedRelations.subList(0, limit); + } else { + return relations; + } } return relationDao.findByRelationPathQuery(tenantId, relationPathQuery, limit); } diff --git a/ui-ngx/src/app/core/auth/auth.models.ts b/ui-ngx/src/app/core/auth/auth.models.ts index d94f26bd0f..623ceefadd 100644 --- a/ui-ngx/src/app/core/auth/auth.models.ts +++ b/ui-ngx/src/app/core/auth/auth.models.ts @@ -35,6 +35,7 @@ export interface SysParamsState { minAllowedAggregationIntervalInSecForCF: number; minAllowedScheduledUpdateIntervalInSecForCF: number; maxRelationLevelPerCfArgument: number; + maxRelatedEntitiesToReturnPerCfArgument: number; ruleChainDebugPerTenantLimitsConfiguration?: string; calculatedFieldDebugPerTenantLimitsConfiguration?: string; intermediateAggregationIntervalInSecForCF: number; diff --git a/ui-ngx/src/app/core/auth/auth.reducer.ts b/ui-ngx/src/app/core/auth/auth.reducer.ts index 41af25f732..51e02b1ab8 100644 --- a/ui-ngx/src/app/core/auth/auth.reducer.ts +++ b/ui-ngx/src/app/core/auth/auth.reducer.ts @@ -37,6 +37,7 @@ const emptyUserAuthState: AuthPayload = { minAllowedAggregationIntervalInSecForCF: 0, minAllowedScheduledUpdateIntervalInSecForCF: 0, maxRelationLevelPerCfArgument: 0, + maxRelatedEntitiesToReturnPerCfArgument: 0, maxDataPointsPerRollingArg: 0, maxDebugModeDurationMinutes: 0, intermediateAggregationIntervalInSecForCF: 0, diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html index 86ee81540b..475f956b07 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html @@ -17,7 +17,7 @@ -->
-
+
{{ 'calculated-fields.propagation-path-related-entities' | translate }}
diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.ts b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.ts index 7aef195480..4eac964538 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.ts +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.ts @@ -43,6 +43,9 @@ import { map } from 'rxjs/operators'; import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; import { ScriptLanguage } from '@app/shared/models/rule-node.models'; import { EntitySearchDirection } from '@shared/models/relation.models'; +import {Store} from "@ngrx/store"; +import {AppState} from "@core/core.state"; +import {getCurrentAuthState} from "@core/auth/auth.selectors"; @Component({ selector: 'tb-propagation-configuration', @@ -79,6 +82,8 @@ export class PropagationConfigurationComponent implements ControlValueAccessor, @Input({transform: booleanAttribute}) isEditValue = true; + readonly maxRelatedEntitiesToReturnPerCfArgument = getCurrentAuthState(this.store).maxRelatedEntitiesToReturnPerCfArgument; + propagateConfiguration = this.fb.group({ arguments: this.fb.control({}, notEmptyObjectValidator()), applyExpressionToResolvedArguments: [false], @@ -112,7 +117,8 @@ export class PropagationConfigurationComponent implements ControlValueAccessor, private propagateChange: (config: CalculatedFieldPropagationConfiguration) => void = () => { }; - constructor(private fb: FormBuilder) { + constructor(private fb: FormBuilder, + private store: Store) { this.propagateConfiguration.get('applyExpressionToResolvedArguments').valueChanges.pipe( takeUntilDestroyed() ).subscribe(() => { diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html b/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html index b0080bb88a..f3ac0eefd8 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html @@ -17,7 +17,7 @@ -->
-
+
{{ 'calculated-fields.aggregation-path-related-entities' | translate }}
diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts b/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts index c2dc66301b..0e21263765 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts @@ -84,6 +84,7 @@ export class RelatedEntitiesAggregationComponentComponent implements ControlValu readonly Directions = Object.values(EntitySearchDirection) as Array; readonly PropagationDirectionTranslations = PropagationDirectionTranslations; readonly minAllowedDeduplicationIntervalInSecForCF = getCurrentAuthState(this.store).minAllowedDeduplicationIntervalInSecForCF; + readonly maxRelatedEntitiesToReturnPerCfArgument = getCurrentAuthState(this.store).maxRelatedEntitiesToReturnPerCfArgument; relatedAggregationConfiguration = this.fb.group({ relation: this.fb.group({ diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index e4b2353cf9..8c9720828a 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1379,9 +1379,9 @@ "zone-group-refresh-interval": "Defines how often zone groups configured via related entities are refreshed.", "zone-group-refresh-interval-required": "Zone groups refresh interval is required.", "zone-group-refresh-interval-min": "Zone group refresh interval should be at least {{ min }} seconds.", - "propagation-path-related-entities": "Defines a direct, single-level path to a related entity based on the selected direction and relation type.", + "propagation-path-related-entities": "Defines a direct, single-level path to a related entity based on the selected direction and relation type. Only relations between device, asset, customer, and tenant entities are supported. Maximum entities resolved by relation path is {{ max }}.", "data-propagate": "Defines the data to be propagated from the arguments configured below. 'Arguments only' uses the retrieved data directly, while 'Expression result' calculates a new value from that data.", - "aggregation-path-related-entities": "Defines a single-level aggregation path via direct relations with parent or child entities, based on direction and relation type. Only relations between device, asset, customer, and tenant entities are supported.", + "aggregation-path-related-entities": "Defines a single-level aggregation path via direct relations with parent or child entities, based on direction and relation type. Only relations between device, asset, customer, and tenant entities are supported. Maximum entities resolved by relation path is {{ max }}.", "arguments-aggregation": "Defines the input arguments used for filtering and aggregation.", "setting-arguments-aggregation": "Data will be fetched from related entities configured in aggregation path.", "metrics": "Defines metrics aggregated based on configured arguments.",