From 7b9b15d0f2a27fd91eaa06243586b4a3d679c708 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 17 Nov 2025 17:06:21 +0200 Subject: [PATCH] review fixes --- .../server/actors/app/AppActor.java | 19 ++++++++++--------- ...CalculatedFieldEntityMessageProcessor.java | 13 +++++-------- ...alculatedFieldManagerMessageProcessor.java | 5 +---- .../EntityInitCalculatedFieldMsg.java | 3 +-- .../cf/ctx/state/CalculatedFieldCtx.java | 16 +++++++++------- ...EntityAggregationCalculatedFieldState.java | 6 +++--- .../alarm/AlarmCalculatedFieldState.java | 3 --- .../src/main/resources/thingsboard.yml | 2 +- 8 files changed, 30 insertions(+), 37 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 5a2f09f789..da16e55db8 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 @@ -166,16 +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 (systemContext.isTenantComponentsInitEnabled()) { + if (msg.getEntityId() instanceof TenantProfileId tenantProfileId) { + tenantService.findTenantIdsByTenantProfileId(tenantProfileId).forEach(tenantId -> { + getOrCreateTenantActor(tenantId).ifPresentOrElse(tenantActor -> { + log.debug("[{}] Sending component lifecycle msg for tenant.", tenantId); + tenantActor.tellWithHighPriority(msg); + }, () -> { + log.debug("Ignoring component lifecycle msg for tenant {} because it is not managed by this service", tenantId); + }); }); - 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/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index b0476b0e34..bf9a3529e5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -123,7 +123,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (state != null) { state.setCtx(msg.getCtx(), actorCtx); state.setPartition(msg.getPartition()); - state.init(true); states.put(cfId, state); } else { removeState(cfId); @@ -134,7 +133,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM log.debug("Processing CF state partition restore msg: {}", msg); for (CalculatedFieldState state : states.values()) { if (msg.getPartition().equals(state.getPartition())) { - state.init(false); + state.init(true); } } } @@ -159,12 +158,10 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } else { state.setCtx(ctx, actorCtx); } - if (msg.getStateAction() != StateAction.REFRESH_CTX) { - if (state.isSizeOk()) { - processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); - } else { - throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); - } + if (state.isSizeOk()) { + processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); + } else { + throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); } } catch (Exception e) { log.debug("[{}][{}] Failed to initialize CF state", entityId, ctx.getCfId(), e); 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 75cf2f6748..ef0abff4dc 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 @@ -260,10 +260,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware 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.REFRESH_CTX, cb)); - }); + callback.onSuccess(); } private void onEntityCreated(ComponentLifecycleMsg msg, TbCallback callback) { diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java index 49f2c691d3..1e0025988d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java @@ -39,7 +39,6 @@ public class EntityInitCalculatedFieldMsg implements ToCalculatedFieldSystemMsg INIT, REINIT, RECREATE, - REPROCESS, - REFRESH_CTX + REPROCESS } } 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 8818ae542c..1d2e2f505a 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 @@ -221,16 +221,18 @@ public class CalculatedFieldCtx implements Closeable { return true; } } - long reevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getAlarmRulesReevaluationInterval()); boolean requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation(); - if (requiresScheduledReevaluation) { - if (now + cfCheckIntervalMillis >= lastReevaluationTs + reevaluationIntervalMillis) { - lastReevaluationTs = now; - return true; + if (calculatedField.getConfiguration() instanceof AlarmCalculatedFieldConfiguration) { + long reevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getAlarmRulesReevaluationInterval()); + if (requiresScheduledReevaluation) { + if (now + cfCheckIntervalMillis >= lastReevaluationTs + reevaluationIntervalMillis) { + lastReevaluationTs = now; + return true; + } + return false; } - return false; } - return false; + return requiresScheduledReevaluation; } public void init() { 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 fbfdd9a258..b07600695c 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 @@ -127,7 +127,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt }); } - public void fillMissingIntervals() { + private void fillMissingIntervals() { ZoneId zoneId = interval.getZoneId(); long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); @@ -253,12 +253,12 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt metricsNode.put(metricName, JacksonUtil.toString(resultValue)); } } - ObjectNode resultNode = JacksonUtil.newObjectNode(); if (!metricsNode.isEmpty()) { + ObjectNode resultNode = JacksonUtil.newObjectNode(); resultNode.put("ts", interval.getEndTs() - 1); resultNode.set("values", metricsNode); + result.add(resultNode); } - result.add(resultNode); }); return result; } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java index 83e08ea67f..342f7534c2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java @@ -124,9 +124,6 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState { @Override public void init(boolean restored) { super.init(restored); - if (restored) { - return; - } AtomicBoolean reevalNeeded = new AtomicBoolean(false); Map createRules = configuration.getCreateRules(); for (AlarmSeverity severity : AlarmSeverity.values()) { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index cdeb55b7b7..b1d66b2f25 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -529,7 +529,7 @@ actors: configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}" # Time in seconds to receive calculation result. calculation_timeout: "${ACTORS_CALCULATION_TIMEOUT_SEC:5}" - # Interval in seconds to re-evaluate calculated fields that have a time schedule. 1 minute by default. + # Interval in seconds to check calculated fields for re-evaluation interval. 1 minute by default. check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:60}" alarms: # Interval in seconds to re-evaluate Alarm rules that have a time schedule. 2 minutes by default.