From 7c886840ce2734da8e9ab4388a72dbe152a6888d Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 27 Oct 2025 17:13:23 +0200 Subject: [PATCH] refactoring --- .../CalculatedFieldEntityActor.java | 4 +- ...CalculatedFieldEntityMessageProcessor.java | 75 ++++++++----------- ...alculatedFieldManagerMessageProcessor.java | 4 +- ... => CalculatedFieldRelationActionMsg.java} | 12 +-- ...tractCalculatedFieldProcessingService.java | 4 +- .../cf/DefaultCalculatedFieldCache.java | 15 +--- .../cf/TelemetryCalculatedFieldResult.java | 2 +- .../server/common/msg/MsgType.java | 2 +- .../thingsboard/script/api/tbel/TbUtils.java | 3 +- 9 files changed, 51 insertions(+), 70 deletions(-) rename application/src/main/java/org/thingsboard/server/actors/calculatedField/{CalculatedFieldRelatedEntityMsg.java => CalculatedFieldRelationActionMsg.java} (78%) diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java index ed4131c114..160cd995d5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java @@ -73,8 +73,8 @@ public class CalculatedFieldEntityActor extends AbstractCalculatedFieldActor { case CF_ENTITY_DELETE_MSG: processor.process((CalculatedFieldEntityDeleteMsg) msg); break; - case CF_RELATED_ENTITY_MSG: - processor.process((CalculatedFieldRelatedEntityMsg) msg); + case CF_RELATION_ACTION_MSG: + processor.process((CalculatedFieldRelationActionMsg) msg); break; case CF_ENTITY_TELEMETRY_MSG: processor.process((EntityCalculatedFieldTelemetryMsg) 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 9cee93611e..ad896f1ef1 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 @@ -210,7 +210,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } - public void process(CalculatedFieldRelatedEntityMsg msg) throws CalculatedFieldException { + public void process(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { log.debug("[{}] Processing CF {} related entity msg.", msg.getRelatedEntityId(), msg.getAction()); switch (msg.getAction()) { case UPDATED -> handleRelationUpdate(msg); @@ -219,7 +219,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } - private void handleRelationUpdate(CalculatedFieldRelatedEntityMsg msg) throws CalculatedFieldException { + private void handleRelationUpdate(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { CalculatedFieldCtx ctx = msg.getCalculatedField(); var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback()); var state = states.get(ctx.getCfId()); @@ -249,7 +249,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } - private void handleRelationDelete(CalculatedFieldRelatedEntityMsg msg) throws CalculatedFieldException { + private void handleRelationDelete(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { CalculatedFieldCtx ctx = msg.getCalculatedField(); CalculatedFieldId cfId = ctx.getCfId(); CalculatedFieldState state = states.get(cfId); @@ -268,7 +268,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); } } else { - // todo: log msg.getCallback().onSuccess(); } } @@ -538,19 +537,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM private Map mapToArguments(EntityId originator, Map argNames, Map relatedEntityArgs, List data) { Map arguments = new HashMap<>(); - if (!relatedEntityArgs.isEmpty()) { + if (!relatedEntityArgs.isEmpty() || !argNames.isEmpty()) { for (TsKvProto item : data) { ReferencedEntityKey key = new ReferencedEntityKey(item.getKv().getKey(), ArgumentType.TS_LATEST, null); String argName = relatedEntityArgs.get(key); if (argName != null) { arguments.put(argName, new SingleValueArgumentEntry(originator, item)); } - } - } - if (!argNames.isEmpty()) { - for (TsKvProto item : data) { - ReferencedEntityKey key = new ReferencedEntityKey(item.getKv().getKey(), ArgumentType.TS_LATEST, null); - String argName = argNames.get(key); + argName = argNames.get(key); if (argName != null) { arguments.put(argName, new SingleValueArgumentEntry(item)); } @@ -577,10 +571,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM private Map mapToArguments(EntityId entityId, Map argNames, List geofencingArgNames, Map relatedEntityArgs, AttributeScopeProto scope, List attrDataList) { Map arguments = new HashMap<>(); - if (!argNames.isEmpty()) { + if (!relatedEntityArgs.isEmpty() || !argNames.isEmpty()) { for (AttributeValueProto item : attrDataList) { ReferencedEntityKey key = new ReferencedEntityKey(item.getKey(), ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); - String argName = argNames.get(key); + String argName = relatedEntityArgs.get(key); + if (argName != null) { + arguments.put(argName, new SingleValueArgumentEntry(entityId, item)); + } + argName = argNames.get(key); if (argName == null) { continue; } @@ -591,15 +589,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM arguments.put(argName, new SingleValueArgumentEntry(item)); } } - if (!relatedEntityArgs.isEmpty()) { - for (AttributeValueProto item : attrDataList) { - ReferencedEntityKey key = new ReferencedEntityKey(item.getKey(), ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); - String argName = relatedEntityArgs.get(key); - if (argName != null) { - arguments.put(argName, new SingleValueArgumentEntry(entityId, item)); - } - } - } return arguments; } @@ -625,23 +614,16 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM AttributeScopeProto scope, List removedAttrKeys) { Map arguments = new HashMap<>(); - if (!relatedEntityArgs.isEmpty()) { - for (String removedKey : removedAttrKeys) { - ReferencedEntityKey key = new ReferencedEntityKey(removedKey, ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); - if (relatedEntityArgs.containsKey(key)) { - String argName = relatedEntityArgs.get(key); - Argument argument = configArguments.get(argName); - String defaultValue = (argument != null) ? argument.getDefaultValue() : null; - SingleValueArgumentEntry argumentEntry = StringUtils.isNotEmpty(defaultValue) - ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null) - : new SingleValueArgumentEntry(); - arguments.put(argName, new SingleValueArgumentEntry(msgEntityId, argumentEntry)); - } - } - } for (String removedKey : removedAttrKeys) { ReferencedEntityKey key = new ReferencedEntityKey(removedKey, ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); - String argName = argNames.get(key); + String argName = relatedEntityArgs.get(key); + if (argName != null) { + String defaultValue = getDefaultValue(configArguments, argName); + SingleValueArgumentEntry argumentEntry = buildSingleValue(removedKey, defaultValue, System.currentTimeMillis()); + arguments.put(argName, new SingleValueArgumentEntry(msgEntityId, argumentEntry)); + continue; + } + argName = argNames.get(key); if (argName == null) { continue; } @@ -649,16 +631,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM arguments.put(argName, new GeofencingArgumentEntry()); continue; } - Argument argument = configArguments.get(argName); - String defaultValue = (argument != null) ? argument.getDefaultValue() : null; - SingleValueArgumentEntry argumentEntry = StringUtils.isNotEmpty(defaultValue) - ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null) - : new SingleValueArgumentEntry(); - arguments.put(argName, argumentEntry); + String defaultValue = getDefaultValue(configArguments, argName); + arguments.put(argName, buildSingleValue(removedKey, defaultValue, System.currentTimeMillis())); } return arguments; } + private String getDefaultValue(Map configArguments, String argName) { + Argument argument = configArguments.get(argName); + return argument != null ? argument.getDefaultValue() : null; + } + + private SingleValueArgumentEntry buildSingleValue(String attrKey, String defaultValue, long ts) { + return StringUtils.isNotEmpty(defaultValue) + ? new SingleValueArgumentEntry(ts, new StringDataEntry(attrKey, defaultValue), null) + : new SingleValueArgumentEntry(); + } + private Map mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, EntityId entityId, List removedTelemetryKeys) { Map deletedArguments = ctx.getArguments().entrySet().stream() .filter(entry -> removedTelemetryKeys.contains(entry.getValue().getRefEntityKey().getKey())) 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 de3967d5b6..40348f8f06 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 @@ -688,12 +688,12 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware private void deleteRelatedEntity(EntityId entityId, EntityId relatedEntityId, CalculatedFieldCtx cf, TbCallback callback) { log.debug("Pushing delete related entity msg to specific actor [{}]", relatedEntityId); - getOrCreateActor(entityId).tell(new CalculatedFieldRelatedEntityMsg(tenantId, relatedEntityId, ActionType.DELETED, cf, callback)); + getOrCreateActor(entityId).tell(new CalculatedFieldRelationActionMsg(tenantId, relatedEntityId, ActionType.DELETED, cf, callback)); } private void initRelatedEntity(EntityId entityId, EntityId relatedEntityId, CalculatedFieldCtx cf, TbCallback callback) { log.debug("Pushing init related entity msg to specific actor [{}]", relatedEntityId); - getOrCreateActor(entityId).tell(new CalculatedFieldRelatedEntityMsg(tenantId, relatedEntityId, ActionType.UPDATED, cf, callback)); + getOrCreateActor(entityId).tell(new CalculatedFieldRelationActionMsg(tenantId, relatedEntityId, ActionType.UPDATED, cf, callback)); } private void deleteCfForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) { diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelatedEntityMsg.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelationActionMsg.java similarity index 78% rename from application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelatedEntityMsg.java rename to application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelationActionMsg.java index bf73a1b1b0..4d8e1cf561 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelatedEntityMsg.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelationActionMsg.java @@ -25,7 +25,7 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; @Data -public class CalculatedFieldRelatedEntityMsg implements ToCalculatedFieldSystemMsg { +public class CalculatedFieldRelationActionMsg implements ToCalculatedFieldSystemMsg { private final TenantId tenantId; private final EntityId relatedEntityId; @@ -33,10 +33,10 @@ public class CalculatedFieldRelatedEntityMsg implements ToCalculatedFieldSystemM private final CalculatedFieldCtx calculatedField; private final TbCallback callback; - public CalculatedFieldRelatedEntityMsg(TenantId tenantId, - EntityId relatedEntityId, ActionType action, - CalculatedFieldCtx calculatedField, - TbCallback callback) { + public CalculatedFieldRelationActionMsg(TenantId tenantId, + EntityId relatedEntityId, ActionType action, + CalculatedFieldCtx calculatedField, + TbCallback callback) { this.tenantId = tenantId; this.relatedEntityId = relatedEntityId; this.action = action; @@ -46,7 +46,7 @@ public class CalculatedFieldRelatedEntityMsg implements ToCalculatedFieldSystemM @Override public MsgType getMsgType() { - return MsgType.CF_RELATED_ENTITY_MSG; + return MsgType.CF_RELATION_ACTION_MSG; } } 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 d4bf0d52be..945792ebcd 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 @@ -49,7 +49,7 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; -import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -188,7 +188,7 @@ public abstract class AbstractCalculatedFieldProcessingService { return Futures.transform(relationsFut, relations -> { if (relations == null) { - return new ArrayList<>(); + return Collections.emptyList(); } return switch (relation.direction()) { 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 1c755ee05a..fd38d3838c 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 @@ -26,8 +26,8 @@ import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; @@ -69,7 +69,6 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { private final ConcurrentMap> calculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); - private final ConcurrentMap aggCalculatedFields = new ConcurrentHashMap<>(); private final ConcurrentMap> ownerEntities = new ConcurrentHashMap<>(); @@ -83,9 +82,6 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { cfs.forEach(cf -> { if (cf != null) { calculatedFields.putIfAbsent(cf.getId(), cf); - if (cf.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration) { - aggCalculatedFields.put(cf.getId(), cf); - } } }); calculatedFields.values().forEach(cf -> { @@ -153,8 +149,8 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { @Override public List getAggCalculatedFieldCtxsByFilter(Predicate relatedEntityFilter) { - return aggCalculatedFields.keySet().stream() - .map(this::getCalculatedFieldCtx) + return calculatedFieldsCtx.values().stream() + .filter(ctx -> CalculatedFieldType.RELATED_ENTITIES_AGGREGATION.equals(ctx.getCfType())) .filter(relatedEntityFilter) .toList(); } @@ -200,9 +196,6 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { entityIdCalculatedFields.computeIfAbsent(cfEntityId, entityId -> new CopyOnWriteArrayList<>()).add(calculatedField); CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); - if (configuration instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration) { - aggCalculatedFields.put(calculatedField.getId(), calculatedField); - } calculatedFieldLinks.put(calculatedFieldId, configuration.buildCalculatedFieldLinks(tenantId, cfEntityId, calculatedFieldId)); configuration.getReferencedEntities().stream() @@ -234,8 +227,6 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { log.debug("[{}] evict calculated field ctx from cache: {}", calculatedFieldId, oldCalculatedField); entityIdCalculatedFieldLinks.forEach((entityId, calculatedFieldLinks) -> calculatedFieldLinks.removeIf(link -> link.getCalculatedFieldId().equals(calculatedFieldId))); log.debug("[{}] evict calculated field links from cached links by entity id: {}", calculatedFieldId, oldCalculatedField); - aggCalculatedFields.remove(calculatedFieldId); - log.debug("[{}] evict calculated field from cached triggers: {}", calculatedFieldId, oldCalculatedField); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java index d59ec9cca9..e71e381807 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java @@ -39,7 +39,7 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu private final AttributeScope scope; private final JsonNode result; - public static TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build(); + public static final TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build(); @Override public TbMsg toTbMsg(EntityId entityId, List cfIds) { diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java index f82de35819..85ce75c829 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java @@ -152,7 +152,7 @@ public enum MsgType { CF_ENTITY_INIT_CF_MSG, CF_ENTITY_DELETE_MSG, - CF_RELATED_ENTITY_MSG, + CF_RELATION_ACTION_MSG, CF_ARGUMENT_RESET_MSG, // Sent to reset argument; CF_REEVALUATE_MSG; diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java index 3e677cf269..39a32310a6 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java @@ -264,6 +264,8 @@ public class TbUtils { float.class, int.class))); parserConfig.addImport("toInt", new MethodStub(TbUtils.class.getMethod("toInt", double.class))); + parserConfig.addImport("roundResult", new MethodStub(TbUtils.class.getMethod("roundResult", + double.class, Integer.class))); parserConfig.addImport("isNaN", new MethodStub(TbUtils.class.getMethod("isNaN", double.class))); parserConfig.addImport("hexToBytes", new MethodStub(TbUtils.class.getMethod("hexToBytes", @@ -1186,7 +1188,6 @@ public class TbUtils { return BigDecimal.valueOf(value).setScale(0, RoundingMode.HALF_UP).intValue(); } - // todo: register method public static Object roundResult(double value, Integer precision) { if (precision == null) { return value;