diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/AbstractTbEntityService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/AbstractTbEntityService.java index d2b3890c70..cef32b9065 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/AbstractTbEntityService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/AbstractTbEntityService.java @@ -89,18 +89,25 @@ public abstract class AbstractTbEntityService { @Lazy private EntitiesVersionControlService vcService; @Autowired + @Lazy protected AccessControlService accessControlService; @Autowired + @Lazy protected TenantService tenantService; @Autowired + @Lazy protected AssetService assetService; @Autowired + @Lazy protected DeviceService deviceService; @Autowired + @Lazy protected AssetProfileService assetProfileService; @Autowired + @Lazy protected DeviceProfileService deviceProfileService; @Autowired + @Lazy protected EntityService entityService; protected boolean isTestProfile() { diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java index 61f28304f2..0a7a5c8264 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java @@ -15,17 +15,20 @@ */ package org.thingsboard.server.service.entitiy.cf; -import lombok.Builder; import lombok.Data; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; @Data -@Builder public class CalculatedFieldCtx { - private final CalculatedFieldId calculatedFieldId; - private final EntityId entityId; - private final CalculatedFieldState state; + private CalculatedFieldId calculatedFieldId; + private EntityId entityId; + private CalculatedFieldState state; + public CalculatedFieldCtx(CalculatedFieldId calculatedFieldId, EntityId entityId, CalculatedFieldState state) { + this.calculatedFieldId = calculatedFieldId; + this.entityId = entityId; + this.state = state; + } } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java index e222259a53..b439d8e641 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.entitiy.cf; +import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; @@ -38,7 +39,6 @@ import org.thingsboard.server.common.data.cf.BaseCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; -import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; @@ -48,8 +48,7 @@ import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.dao.attributes.AttributesService; @@ -62,6 +61,7 @@ import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; @@ -71,7 +71,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; import static org.thingsboard.server.dao.service.Validator.validateEntityId; @@ -128,98 +127,6 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp } } - private ListenableFuture> fetchAttributesForEntity(TenantId tenantId, EntityId entityId, List keys) { - return attributesService.find(tenantId, entityId, AttributeScope.SERVER_SCOPE, keys); - } - - private ListenableFuture> fetchTimeSeries(TenantId tenantId, EntityId entityId, List keys) { - return timeseriesService.findLatest(tenantId, entityId, keys); - } - - private ListenableFuture initializeStateFromFutures(TenantId tenantId, EntityId entityId, CalculatedField calculatedField, List attributeKeys, List timeSeriesKeys) { - ListenableFuture> attributesFuture = fetchAttributesForEntity(tenantId, entityId, attributeKeys); - ListenableFuture> timeSeriesFuture = fetchTimeSeries(tenantId, entityId, timeSeriesKeys); - - ListenableFuture> combinedFuture = Futures.allAsList(attributesFuture, timeSeriesFuture); - - return Futures.transform(combinedFuture, results -> { - List attributes = (List) results.get(0); - List timeSeries = (List) results.get(1); - - initializeState(calculatedField, attributes, timeSeries); - - return null; - }, MoreExecutors.directExecutor()); - } - - private void initializeState(CalculatedField calculatedField, List attributes, List timeSeries) { - CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(calculatedField.getId(), - ctx -> new CalculatedFieldCtx(calculatedField.getId(), calculatedField.getEntityId(), null)); - - CalculatedFieldState state = calculatedFieldCtx.getState(); - - if (state != null) { - String calculation = performCalculation(state.getArguments()); - - Map updatedArguments = state.getArguments(); - - state = CalculatedFieldState.builder() - .arguments(updatedArguments) - .result(calculation) - .build(); - } else { - // initial calculation - Map arguments = calculatedField.getConfiguration().getArguments(); - - Map argumentValues = arguments.entrySet().stream() - .collect(Collectors.toMap( - Map.Entry::getKey, - entry -> resolveArgumentValue(entry.getKey(), entry.getValue(), attributes, timeSeries) - )); - - String calculation = performCalculation(argumentValues); - - state = CalculatedFieldState.builder() - .arguments(argumentValues) - .result(calculation) - .build(); - } - - calculatedFieldCtx = new CalculatedFieldCtx(calculatedField.getId(), calculatedField.getEntityId(), state); - states.put(calculatedField.getId(), calculatedFieldCtx); - } - - private String resolveArgumentValue(String key, BaseCalculatedFieldConfiguration.Argument argument, - List attributes, List timeSeries) { - String type = argument.getType(); - String value = null; - - if ("ATTRIBUTES".equals(type)) { - value = attributes.stream() - .filter(attribute -> attribute.getKey().equals(key)) - .map(AttributeKvEntry::getValueAsString) - .findFirst() - .orElse(null); - } else if ("TIME_SERIES".equals(type)) { - value = timeSeries.stream() - .filter(tsKvEntry -> tsKvEntry.getKey().equals(key)) - .map(TsKvEntry::getValueAsString) - .findFirst() - .orElse(null); - } - - return value != null ? value : argument.getDefaultValue(); - } - - private void initializeForProfile(TenantId tenantId, List links, CalculatedField cf, Iterable profileIds) { - for (T profileId : profileIds) { - for (CalculatedFieldLink link : links) { - CalculatedFieldLinkConfiguration configuration = link.getConfiguration(); - initializeStateFromFutures(tenantId, profileId, cf, configuration.getAttributes(), configuration.getTimeSeries()); - } - } - } - @Override public void onCalculatedFieldAdded(TransportProtos.CalculatedFieldAddMsgProto proto, TbCallback callback) { try { @@ -232,21 +139,16 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp calculatedFields.put(calculatedFieldId, cf); calculatedFieldLinks.put(calculatedFieldId, links); switch (entityId.getEntityType()) { - case ASSET, DEVICE -> { - for (CalculatedFieldLink link : links) { - CalculatedFieldLinkConfiguration configuration = link.getConfiguration(); - initializeStateFromFutures(tenantId, link.getEntityId(), cf, configuration.getAttributes(), configuration.getTimeSeries()); - } - } + case ASSET, DEVICE -> initializeStateForEntity(tenantId, cf, callback); case ASSET_PROFILE -> { PageDataIterable assetIds = new PageDataIterable<>(pageLink -> assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) entityId, pageLink), initFetchPackSize); - initializeForProfile(tenantId, links, cf, assetIds); + assetIds.forEach(assetId -> initializeStateForEntity(tenantId, cf, callback)); } case DEVICE_PROFILE -> { PageDataIterable deviceIds = new PageDataIterable<>(pageLink -> deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityId, pageLink), initFetchPackSize); - initializeForProfile(tenantId, links, cf, deviceIds); + deviceIds.forEach(deviceId -> initializeStateForEntity(tenantId, cf, callback)); } default -> throw new IllegalArgumentException("Entity type '" + calculatedFieldId.getEntityType() + "' does not support calculated fields."); @@ -359,7 +261,72 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp }; } - private String performCalculation(Map argumentValues) { + private void initializeStateForEntity(TenantId tenantId, CalculatedField calculatedField, TbCallback callback) { + Map arguments = calculatedField.getConfiguration().getArguments(); + Map argumentValues = new HashMap<>(); + arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() { + @Override + public void onSuccess(Optional result) { + String value = result.map(KvEntry::getValueAsString).orElse(argument.getDefaultValue()); + argumentValues.put(key, value); + } + + @Override + public void onFailure(Throwable t) { + log.warn("Failed to fetch data for type: {}", argument.getType(), t); + callback.onFailure(t); + } + }, calculatedFieldCallbackExecutor)); + + updateOrInitializeState(calculatedField, argumentValues); + + } + + private ListenableFuture> fetchArgumentValue(TenantId tenantId, BaseCalculatedFieldConfiguration.Argument argument) { + return switch (argument.getType()) { + case "ATTRIBUTES" -> Futures.transform( + attributesService.find(tenantId, argument.getEntityId(), AttributeScope.SERVER_SCOPE, argument.getKey()), + result -> result.map(entry -> (KvEntry) entry), + MoreExecutors.directExecutor()); + case "TIME_SERIES" -> Futures.transform( + timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()), + result -> result.map(entry -> (KvEntry) entry), + MoreExecutors.directExecutor()); + default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); + }; + } + + private void updateOrInitializeState(CalculatedField calculatedField, Map argumentValues) { + CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(calculatedField.getId(), + ctx -> new CalculatedFieldCtx(calculatedField.getId(), calculatedField.getEntityId(), null)); + + CalculatedFieldState state = calculatedFieldCtx.getState(); + + if (state != null) { + // calculation based on the previous data + String calculation = performCalculation(state.getArguments(), calculatedField.getConfiguration()); + + Map updatedArguments = new HashMap<>(state.getArguments()); + + state = CalculatedFieldState.builder() + .arguments(updatedArguments) + .result(calculation) + .build(); + } else { + // initial calculation + String calculation = performCalculation(argumentValues, calculatedField.getConfiguration()); + + state = CalculatedFieldState.builder() + .arguments(argumentValues) + .result(calculation) + .build(); + } + calculatedFieldCtx.setState(state); + states.put(calculatedField.getId(), calculatedFieldCtx); + } + + private String performCalculation(Map argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration) { + BaseCalculatedFieldConfiguration.Output output = calculatedFieldConfiguration.getOutput(); return "calculation"; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java index deaeebcb50..cfae2e5a0e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java @@ -41,6 +41,9 @@ public interface CalculatedFieldConfiguration { Map getArguments(); + @JsonIgnore + BaseCalculatedFieldConfiguration.Output getOutput(); + @JsonIgnore List getReferencedEntities();