Browse Source

Merge pull request #14770 from irynamatveieva/fix/related-entities-cf

Related entities aggregation CF: fixes
pull/14785/head
Viacheslav Klimov 9 months ago
committed by GitHub
parent
commit
320583ea3a
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 51
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  2. 1
      application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java
  3. 10
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  4. 28
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  5. 18
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  6. 27
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java
  7. 12
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java
  8. 4
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java
  9. 4
      application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java
  10. 101
      application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java
  11. 2
      application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java
  12. 3
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java
  13. 3
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesArgumentEntryTest.java
  14. 1
      common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java
  15. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java
  16. 19
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  17. 1
      ui-ngx/src/app/core/auth/auth.models.ts
  18. 1
      ui-ngx/src/app/core/auth/auth.reducer.ts
  19. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html
  20. 8
      ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.ts
  21. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html
  22. 1
      ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts
  23. 4
      ui-ngx/src/assets/locale/locale.constant-en_US.json

51
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<CalculatedFieldCtx> 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<EntityId> 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())));
}
}
}

1
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());

10
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);
}

28
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;

18
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() + ".");
}
}
}
}
}

27
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

12
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<EntityId> updatedIds) {
private boolean checkAdded(Collection<EntityId> 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<>();

4
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();
}

4
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();

101
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<String, Argument> arguments = new HashMap<>();
Argument argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("active", ArgumentType.TS_LATEST, null));
argument.setDefaultValue("true");
arguments.put("active", argument);
Map<String, AggMetric> 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,

2
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)));

3
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<EntityId> propagationEntityIds = new ArrayList<>();
propagationEntityIds.add(ENTITY_1_ID);
propagationEntityIds.add(ENTITY_2_ID);

3
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<EntityId, ArgumentEntry> 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)));

1
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;

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

19
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<EntityRelation> relations = relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations;
return relations.size() > limit ? relations.subList(0, limit) : relations;
if (relations.size() > limit) {
List<EntityRelation> 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<EntityRelation> 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);
}

1
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;

1
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,

2
ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html

@ -17,7 +17,7 @@
-->
<div [formGroup]="propagateConfiguration" class="tb-form-panel no-border no-padding">
<div class="tb-form-panel">
<div class="tb-form-panel-title" tbTruncateWithTooltip tb-hint-tooltip-icon="{{ 'calculated-fields.hint.propagation-path-related-entities' | translate }}">
<div class="tb-form-panel-title" tbTruncateWithTooltip tb-hint-tooltip-icon="{{ 'calculated-fields.hint.propagation-path-related-entities' | translate: {max: maxRelatedEntitiesToReturnPerCfArgument} }}">
{{ 'calculated-fields.propagation-path-related-entities' | translate }}
</div>
<div class="flex gap-3 xs:flex-col" formGroupName="relation">

8
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<AppState>) {
this.propagateConfiguration.get('applyExpressionToResolvedArguments').valueChanges.pipe(
takeUntilDestroyed()
).subscribe(() => {

2
ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.html

@ -17,7 +17,7 @@
-->
<div [formGroup]="relatedAggregationConfiguration" class="tb-form-panel no-border no-padding">
<div class="tb-form-panel">
<div class="tb-form-panel-title" tbTruncateWithTooltip tb-hint-tooltip-icon="{{ 'calculated-fields.hint.aggregation-path-related-entities' | translate }}">
<div class="tb-form-panel-title" tbTruncateWithTooltip tb-hint-tooltip-icon="{{ 'calculated-fields.hint.aggregation-path-related-entities' | translate: {max: maxRelatedEntitiesToReturnPerCfArgument} }}">
{{ 'calculated-fields.aggregation-path-related-entities' | translate }}
</div>
<div class="flex gap-3 xs:flex-col" formGroupName="relation">

1
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<EntitySearchDirection>;
readonly PropagationDirectionTranslations = PropagationDirectionTranslations;
readonly minAllowedDeduplicationIntervalInSecForCF = getCurrentAuthState(this.store).minAllowedDeduplicationIntervalInSecForCF;
readonly maxRelatedEntitiesToReturnPerCfArgument = getCurrentAuthState(this.store).maxRelatedEntitiesToReturnPerCfArgument;
relatedAggregationConfiguration = this.fb.group({
relation: this.fb.group({

4
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.",

Loading…
Cancel
Save