From 1fd97498bea25c0192b6bea4b7fa4f5088748fd2 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 13 Nov 2025 10:42:44 +0200 Subject: [PATCH] handle tenant profile update --- .../server/actors/app/AppActor.java | 12 +++++++++ ...alculatedFieldManagerMessageProcessor.java | 19 ++++++++++++++ .../server/actors/tenant/TenantActor.java | 2 +- .../service/cf/CalculatedFieldCache.java | 2 ++ .../cf/DefaultCalculatedFieldCache.java | 5 ++++ .../cf/ctx/state/CalculatedFieldCtx.java | 26 ++++++++++++------- ...EntityAggregationCalculatedFieldState.java | 2 +- .../processing/AbstractConsumerService.java | 1 + .../thingsboard/server/cf/AlarmRulesTest.java | 3 ++- .../EntityAggregationCalculatedFieldTest.java | 5 ++-- 10 files changed, 63 insertions(+), 14 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index 20cacda26a..5a2f09f789 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -32,6 +32,7 @@ import org.thingsboard.server.actors.tenant.TenantActor; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.MsgType; @@ -165,6 +166,17 @@ public class AppActor extends ContextAwareActor { private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) { TbActorRef target = null; if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) { + if (msg.getEntityId() instanceof TenantProfileId tenantProfileId) { + tenantService.findTenantIdsByTenantProfileId(tenantProfileId).forEach(tenantId -> { + TbActorRef tenantActor = getOrCreateTenantActor(tenantId).orElseGet(() -> { + log.debug("Ignoring component lifecycle msg for tenant {} because it is not managed by this service", tenantId); + return null; + }); + if (tenantActor != null) { + tenantActor.tellWithHighPriority(msg); + } + }); + } if (!msg.getEntityId().getEntityType().isOneOf(EntityType.TENANT_PROFILE, EntityType.TB_RESOURCE)) { log.warn("Message has system tenant id: {}", msg); } 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 a5b975d454..f0036b6627 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 @@ -71,6 +71,7 @@ import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; import java.util.ArrayList; +import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -83,6 +84,7 @@ import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; import java.util.function.Function; +import java.util.stream.Stream; import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; @@ -222,6 +224,12 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware default -> msg.getCallback().onSuccess(); } } + case TENANT_PROFILE -> { + switch (event) { + case UPDATED -> onTenantProfileUpdated(msg.getData(), msg.getCallback()); + default -> msg.getCallback().onSuccess(); + } + } default -> msg.getCallback().onSuccess(); } } @@ -247,6 +255,17 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware callback.onSuccess(); } + private void onTenantProfileUpdated(ComponentLifecycleMsg msg, TbCallback callback) { + Stream.concat( + calculatedFields.values().stream(), + entityIdCalculatedFields.values().stream().flatMap(Collection::stream) + ).forEach(CalculatedFieldCtx::updateTenantProfileProperties); + + calculatedFields.values().forEach(ctx -> { + applyToTargetCfEntityActors(ctx, callback, (id, cb) -> initCfForEntity(id, ctx, StateAction.REPROCESS, cb)); + }); + } + private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) { EntityId entityId = msg.getEntityId(); EntityId profileId = getProfileId(tenantId, entityId); diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index 35a7f01b2e..11a8651026 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -350,7 +350,7 @@ public class TenantActor extends RuleChainManagerActor { } } if (cfActor != null) { - if (msg.getEntityId().getEntityType().isOneOf(EntityType.CALCULATED_FIELD, EntityType.DEVICE, EntityType.ASSET, EntityType.CUSTOMER)) { + if (msg.getEntityId().getEntityType().isOneOf(EntityType.CALCULATED_FIELD, EntityType.DEVICE, EntityType.ASSET, EntityType.CUSTOMER, EntityType.TENANT_PROFILE)) { cfActor.tellWithHighPriority(new CalculatedFieldEntityLifecycleMsg(tenantId, msg)); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java index d50a125451..5d643908ce 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java @@ -50,6 +50,8 @@ public interface CalculatedFieldCache { void evict(CalculatedFieldId calculatedFieldId); + void handleTenantProfileUpdate(); + EntityId getProfileId(TenantId tenantId, EntityId entityId); Set getDynamicEntities(TenantId tenantId, EntityId entityId); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java index 9e51997a20..1938e58a8f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -229,6 +229,11 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { log.debug("[{}] evict calculated field links from cached links by entity id: {}", calculatedFieldId, oldCalculatedField); } + @Override + public void handleTenantProfileUpdate() { + calculatedFieldsCtx.values().forEach(CalculatedFieldCtx::updateTenantProfileProperties); + } + @Override public EntityId getProfileId(TenantId tenantId, EntityId entityId) { return switch (entityId.getEntityType()) { 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 518b65dc5f..8818ae542c 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 @@ -97,9 +97,6 @@ public class CalculatedFieldCtx implements Closeable { private String expression; private boolean useLatestTs; - private long cfCheckInterval; - private long alarmReevaluationInterval; - private long lastReevaluationTs; private ActorSystemContext systemContext; @@ -113,7 +110,6 @@ public class CalculatedFieldCtx implements Closeable { private boolean initialized; - private long maxDataPointsPerRollingArg; private long maxStateSize; private long maxSingleValueArgumentSize; @@ -202,15 +198,12 @@ public class CalculatedFieldCtx implements Closeable { if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig) { this.useLatestTs = aggConfig.isUseLatestTs(); } - this.cfCheckInterval = systemContext.getCfCheckInterval(); - this.alarmReevaluationInterval = systemContext.getAlarmRulesReevaluationInterval(); this.systemContext = systemContext; this.tbelInvokeService = systemContext.getTbelInvokeService(); this.relationService = systemContext.getRelationService(); this.alarmService = systemContext.getAlarmService(); this.cfProcessingService = systemContext.getCalculatedFieldProcessingService(); - this.maxDataPointsPerRollingArg = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); // fixme why tenant profile update is not handled?? this.maxStateSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; this.maxSingleValueArgumentSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; } @@ -284,6 +277,11 @@ public class CalculatedFieldCtx implements Closeable { } } + public void updateTenantProfileProperties() { + this.maxStateSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; + this.maxSingleValueArgumentSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; + } + public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) { for (Map.Entry entry : state.getArguments().entrySet()) { try { @@ -645,9 +643,8 @@ public class CalculatedFieldCtx implements Closeable { if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig && other.getCalculatedField().getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) { boolean metricsChanged = thisConfig.getMetrics().equals(otherConfig.getMetrics()); - boolean intervalChanged = thisConfig.getInterval().equals(otherConfig.getInterval()); boolean watermarkChanged = thisConfig.getWatermark().equals(otherConfig.getWatermark()); - return metricsChanged || intervalChanged || watermarkChanged; + return metricsChanged || watermarkChanged; } return false; } @@ -672,6 +669,9 @@ public class CalculatedFieldCtx implements Closeable { if (hasRelatedEntitiesAggregationConfigurationChanges(other)) { return true; } + if (hasEntityAggregationConfigurationChanges(other)) { + return true; + } return false; } @@ -691,6 +691,14 @@ public class CalculatedFieldCtx implements Closeable { return false; } + private boolean hasEntityAggregationConfigurationChanges(CalculatedFieldCtx other) { + if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig + && other.calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) { + return !thisConfig.getInterval().equals(otherConfig.getInterval()); + } + return false; + } + private boolean isScheduledUpdateEnabled() { return scheduledUpdateIntervalMillis != -1; } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java index a1208c5054..04179bdd31 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java @@ -69,7 +69,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); Watermark watermark = configuration.getWatermark(); watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); - checkInterval = TimeUnit.SECONDS.toMillis(ctx.getCfCheckInterval()); + checkInterval = TimeUnit.SECONDS.toMillis(ctx.getSystemContext().getCfCheckInterval()); interval = configuration.getInterval(); metrics = configuration.getMetrics(); } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index 6e162256a4..37c3d31d0a 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -166,6 +166,7 @@ public abstract class AbstractConsumerService { tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1); + tenantProfileConfig.setMinAllowedAggregationIntervalInSecForCF(1); }); Tenant tenant = new Tenant(); @@ -95,7 +96,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCf_checkAggregation() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L); + CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 30L); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); @@ -126,7 +127,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest public void testCreateCf_checkAggregationDuringWatermark() throws Exception { Device device = createDevice("Device", "1234567890111"); - CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L); + CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 30L); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs();