Browse Source

handle tenant profile update

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
1fd97498be
  1. 12
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  2. 19
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 2
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java
  5. 5
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  6. 26
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  8. 1
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  9. 3
      application/src/test/java/org/thingsboard/server/cf/AlarmRulesTest.java
  10. 5
      application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java

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

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

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

2
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<EntityId> getDynamicEntities(TenantId tenantId, EntityId entityId);

5
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()) {

26
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<String, ArgumentEntry> 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;
}

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

1
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -166,6 +166,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
tenantProfileCache.evict(tenantProfileId);
if (componentLifecycleMsg.getEvent().equals(ComponentLifecycleEvent.UPDATED)) {
apiUsageStateService.onTenantProfileUpdate(tenantProfileId);
calculatedFieldCache.handleTenantProfileUpdate();
}
} else if (EntityType.TENANT.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
if (TenantId.SYS_TENANT_ID.equals(tenantId)) {

3
application/src/test/java/org/thingsboard/server/cf/AlarmRulesTest.java

@ -86,7 +86,8 @@ import static org.testcontainers.shaded.org.awaitility.Awaitility.await;
@Slf4j
@DaoSqlTest
@TestPropertySource(properties = {
"actors.calculated_fields.check_interval=1"
"actors.calculated_fields.check_interval=1",
"actors.alarms.reevaluation_interval=1"
})
public class AlarmRulesTest extends AbstractControllerTest {

5
application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java

@ -67,6 +67,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest
updateDefaultTenantProfileConfig(tenantProfileConfig -> {
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();

Loading…
Cancel
Save