Browse Source

refactoring

pull/14141/head
IrynaMatveieva 11 months ago
parent
commit
7c886840ce
  1. 4
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java
  2. 75
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  3. 4
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  4. 12
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelationActionMsg.java
  5. 4
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  6. 15
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java
  8. 2
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  9. 3
      common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java

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

75
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<String, ArgumentEntry> mapToArguments(EntityId originator, Map<ReferencedEntityKey, String> argNames, Map<ReferencedEntityKey, String> relatedEntityArgs, List<TsKvProto> data) {
Map<String, ArgumentEntry> 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<String, ArgumentEntry> mapToArguments(EntityId entityId, Map<ReferencedEntityKey, String> argNames, List<String> geofencingArgNames, Map<ReferencedEntityKey, String> relatedEntityArgs, AttributeScopeProto scope, List<AttributeValueProto> attrDataList) {
Map<String, ArgumentEntry> 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<String> removedAttrKeys) {
Map<String, ArgumentEntry> 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<String, Argument> 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<String, ArgumentEntry> mapToArgumentsWithFetchedValue(CalculatedFieldCtx ctx, EntityId entityId, List<String> removedTelemetryKeys) {
Map<String, Argument> deletedArguments = ctx.getArguments().entrySet().stream()
.filter(entry -> removedTelemetryKeys.contains(entry.getValue().getRefEntityKey().getKey()))

4
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) {

12
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldRelatedEntityMsg.java → 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;
}
}

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

15
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<CalculatedFieldId, List<CalculatedFieldLink>> calculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<EntityId, List<CalculatedFieldLink>> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, CalculatedFieldCtx> calculatedFieldsCtx = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, CalculatedField> aggCalculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap<EntityId, Set<EntityId>> 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<CalculatedFieldCtx> getAggCalculatedFieldCtxsByFilter(Predicate<CalculatedFieldCtx> 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

2
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<CalculatedFieldId> cfIds) {

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

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

Loading…
Cancel
Save