From 3ad0b4b74de092e34fb375aa11a99e952dcfa414 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 8 Jan 2026 16:57:05 +0200 Subject: [PATCH 1/4] fixed relation update handling, handle max related entities limit --- ...alculatedFieldManagerMessageProcessor.java | 37 +++++-- ...tractCalculatedFieldProcessingService.java | 8 +- .../cf/ctx/state/CalculatedFieldCtx.java | 8 ++ .../RelatedEntitiesArgumentEntry.java | 12 +++ .../propagation/PropagationArgumentEntry.java | 12 ++- ...ntitiesAggregationCalculatedFieldTest.java | 101 ++++++++++++++++-- .../dao/relation/BaseRelationService.java | 4 +- 7 files changed, 161 insertions(+), 21 deletions(-) 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..f06c0f113a 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); + Set cfsToReinit = new HashSet<>(); + 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 { @@ -581,6 +596,10 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware if (byRelationPathQuery != null && !byRelationPathQuery.isEmpty()) { switch (relation.direction()) { case FROM -> { + if (byRelationPathQuery.size() > 1) { + throw new IllegalStateException("More than one relation found with direction 'TO' " + + "for relation type '" + relation.relationType() + "'. Found: " + byRelationPathQuery.size()); + } EntityRelation entityRelation = byRelationPathQuery.get(0); // only one supported EntityId relatedId = entityRelation.getFrom(); if (matchesCfEntity.test(relatedId)) { 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..8284e81f70 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 @@ -250,11 +250,17 @@ public abstract class AbstractCalculatedFieldProcessingService { case FROM -> relations.stream() .map(EntityRelation::getTo) .toList(); - case TO -> relations.stream() + case TO -> { + if (relations.size() > 1) { + throw new IllegalStateException("More than one relation found with direction 'TO' " + + "for relation type '" + relation.relationType() + "'. Found: " + relations.size()); + } + yield relations.stream() .map(EntityRelation::getFrom) .findFirst() .map(List::of) .orElseGet(Collections::emptyList); + } }; }, 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..f74ec7514d 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 @@ -129,6 +129,7 @@ public class CalculatedFieldCtx implements Closeable { private long scheduledUpdateIntervalMillis; private long cfCheckReevaluationIntervalMillis; private long alarmReevaluationIntervalMillis; + private long maxRelatedEntitiesPerCfArgument; private Argument propagationArgument; private boolean applyExpressionForResolvedArguments; @@ -307,6 +308,7 @@ public class CalculatedFieldCtx implements Closeable { 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)); + this.maxRelatedEntitiesPerCfArgument = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxRelatedEntitiesToReturnPerCfArgument); } public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) { @@ -756,6 +758,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/RelatedEntitiesArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java index 1939318e7a..c1c8d38316 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,10 +66,12 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs @Override public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { + checkMaxRelatedEntitiesPerArgument(ctx); entityInputs.putAll(relatedEntitiesArgumentEntry.entityInputs); return true; } else if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { if (entry.isForceResetPrevious()) { + checkMaxRelatedEntitiesPerArgument(ctx); entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); return true; } @@ -77,6 +79,7 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs if (argumentEntry != null) { argumentEntry.updateEntry(singleValueArgumentEntry, ctx); } else { + checkMaxRelatedEntitiesPerArgument(ctx); entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); } return true; @@ -85,6 +88,15 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs } } + private void checkMaxRelatedEntitiesPerArgument(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 public boolean isEmpty() { return entityInputs.isEmpty(); 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/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/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..b390b732d5 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 @@ -60,6 +60,7 @@ import org.thingsboard.server.dao.usagerecord.ApiLimitService; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -517,7 +518,8 @@ class BaseRelationService implements RelationService { if (entityRelations == null || entityRelations.isEmpty()) { return Collections.emptyList(); } - List relations = relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations; + List relations = new ArrayList<>(relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations); + relations.sort(Comparator.comparing(r -> r.getFrom().getId())); return relations.size() > limit ? relations.subList(0, limit) : relations; }, directExecutor()); } From 7a638c2151c6bc4c1cc7e2ea0b1d11766edfd463 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 9 Jan 2026 14:33:05 +0200 Subject: [PATCH 2/4] updated metric default value to double and added logic to handle relations when direction to --- ...alculatedFieldManagerMessageProcessor.java | 22 +++++------------- ...tractCalculatedFieldProcessingService.java | 16 ++----------- ...titiesAggregationCalculatedFieldState.java | 18 +++++++++++++++ .../RelatedEntitiesArgumentEntry.java | 23 +++++++++---------- .../utils/CalculatedFieldArgumentUtils.java | 4 ++-- .../EntityAggregationCalculatedFieldTest.java | 4 ++-- .../CalculatedFieldControllerTest.java | 2 +- .../configuration/aggregation/AggMetric.java | 2 +- .../dao/relation/BaseRelationService.java | 6 +++-- 9 files changed, 47 insertions(+), 50 deletions(-) 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 f06c0f113a..452435ae9d 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 @@ -595,22 +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 -> { - if (byRelationPathQuery.size() > 1) { - throw new IllegalStateException("More than one relation found with direction 'TO' " + - "for relation type '" + relation.relationType() + "'. Found: " + byRelationPathQuery.size()); - } - 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/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 8284e81f70..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,20 +247,8 @@ public abstract class AbstractCalculatedFieldProcessingService { } return switch (relation.direction()) { - case FROM -> relations.stream() - .map(EntityRelation::getTo) - .toList(); - case TO -> { - if (relations.size() > 1) { - throw new IllegalStateException("More than one relation found with direction 'TO' " + - "for relation type '" + relation.relationType() + "'. Found: " + relations.size()); - } - yield 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/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 c1c8d38316..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,29 +66,28 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs @Override public boolean updateEntry(ArgumentEntry entry, CalculatedFieldCtx ctx) { if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { - checkMaxRelatedEntitiesPerArgument(ctx); + checkRelatedEntitiesNumber(ctx); entityInputs.putAll(relatedEntitiesArgumentEntry.entityInputs); - return true; } else if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { if (entry.isForceResetPrevious()) { - checkMaxRelatedEntitiesPerArgument(ctx); + checkRelatedEntitiesNumber(ctx); entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); - return true; - } - ArgumentEntry argumentEntry = entityInputs.get(singleValueArgumentEntry.getEntityId()); - if (argumentEntry != null) { - argumentEntry.updateEntry(singleValueArgumentEntry, ctx); } else { - checkMaxRelatedEntitiesPerArgument(ctx); - 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 checkMaxRelatedEntitiesPerArgument(CalculatedFieldCtx ctx) { + private void checkRelatedEntitiesNumber(CalculatedFieldCtx ctx) { if (entityInputs.size() >= ctx.getMaxRelatedEntitiesPerCfArgument()) { throw new IllegalArgumentException( "Exceeded the maximum allowed related entities per argument '" 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/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/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 b390b732d5..166fba81ff 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,11 +52,11 @@ 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; @@ -547,7 +547,9 @@ 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; + ArrayList entityRelations = new ArrayList<>(relations); + entityRelations.sort(Comparator.comparing(r -> r.getFrom().getId())); + return entityRelations.size() > limit ? entityRelations.subList(0, limit) : entityRelations; } return relationDao.findByRelationPathQuery(tenantId, relationPathQuery, limit); } From 9813f67d576a3e727ccb3493ceae78f6614bb76a Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 12 Jan 2026 10:11:43 +0200 Subject: [PATCH 3/4] fixed tests and updated hints for relation path --- .../server/controller/SystemInfoController.java | 1 + .../cf/ctx/state/PropagationArgumentEntryTest.java | 3 +++ .../cf/ctx/state/RelatedEntitiesArgumentEntryTest.java | 3 +++ .../org/thingsboard/server/common/data/SystemParams.java | 1 + ui-ngx/src/app/core/auth/auth.models.ts | 1 + ui-ngx/src/app/core/auth/auth.reducer.ts | 1 + .../propagation-configuration.component.html | 2 +- .../propagation-configuration.component.ts | 8 +++++++- .../related-entities-aggregation-component.component.html | 2 +- .../related-entities-aggregation-component.component.ts | 1 + ui-ngx/src/assets/locale/locale.constant-en_US.json | 4 ++-- 11 files changed, 22 insertions(+), 5 deletions(-) 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/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/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.", From 5565b5c32a296d33e8f319a5882bdf1bf8728aba Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 12 Jan 2026 12:11:08 +0200 Subject: [PATCH 4/4] code optimization --- ...alculatedFieldManagerMessageProcessor.java | 2 +- .../cf/ctx/state/CalculatedFieldCtx.java | 22 +++++++++++-------- .../dao/relation/BaseRelationService.java | 21 +++++++++++++----- 3 files changed, 29 insertions(+), 16 deletions(-) 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 452435ae9d..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 @@ -265,7 +265,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware checkCfIntervalForUpdate(); long maxRelatedEntitiesPerCfArgument = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxRelatedEntitiesToReturnPerCfArgument); - Set cfsToReinit = new HashSet<>(); + List cfsToReinit = new ArrayList<>(); Stream.concat( calculatedFields.values().stream(), entityIdCalculatedFields.values().stream().flatMap(Collection::stream) 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 f74ec7514d..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; @@ -302,13 +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)); - this.maxRelatedEntitiesPerCfArgument = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxRelatedEntitiesToReturnPerCfArgument); + 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) { 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 166fba81ff..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 @@ -518,9 +518,14 @@ class BaseRelationService implements RelationService { if (entityRelations == null || entityRelations.isEmpty()) { return Collections.emptyList(); } - List relations = new ArrayList<>(relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations); - relations.sort(Comparator.comparing(r -> r.getFrom().getId())); - return relations.size() > limit ? relations.subList(0, limit) : relations; + List relations = relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations; + 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(() -> { @@ -547,9 +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); }; - ArrayList entityRelations = new ArrayList<>(relations); - entityRelations.sort(Comparator.comparing(r -> r.getFrom().getId())); - return entityRelations.size() > limit ? entityRelations.subList(0, limit) : entityRelations; + 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); }